File size: 6,613 Bytes
7df674a | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 | #!/usr/bin/env python3
"""harvest pod benchmark — throughput + GPU-utilization sweep.
Runs advisory-style chat completions at increasing concurrency against a local
llama-server, with reasoning ON or OFF, while sampling nvidia-smi. Reports
aggregate output tok/s (the budget number), latency, TTFT, and GPU util/VRAM.
python3 bench-pod.py --base-url http://127.0.0.1:18080/v1 --model qwen3.6-35b \
--concurrencies 1,8,32,64 --thinking off --max-tokens 256 --out r.json
"""
import argparse, json, statistics, subprocess, sys, threading, time
from concurrent.futures import ThreadPoolExecutor, as_completed
try:
import requests
except Exception:
sys.exit("needs `requests` (pip install requests)")
SYSTEM = ("You are one advisory agent in a multi-agent system serving a smallholder "
"farmer across a full cropping season. Give concise, decision-focused advice.")
USER = ("Farmer profile: 0.8 ha rainfed plot, sandy-loam soil, monsoon onset delayed ~2 weeks, "
"cotton, moderate bollworm pressure reported nearby, limited cash for inputs. "
"It is the sowing window. Recommend: sow now vs wait, variety duration, and the single "
"most important early-season action. Explain your reasoning in a few sentences.")
class GpuSampler(threading.Thread):
def __init__(self):
super().__init__(daemon=True)
self.samples, self.on = [], True
def run(self):
while self.on:
try:
out = subprocess.run(
["nvidia-smi", "--query-gpu=utilization.gpu,memory.used",
"--format=csv,noheader,nounits"],
capture_output=True, text=True, timeout=5).stdout.strip().splitlines()[0]
u, m = [int(x.strip()) for x in out.split(",")]
self.samples.append((u, m))
except Exception:
pass
time.sleep(1)
def window(self, start_idx):
w = self.samples[start_idx:]
if not w: return {}
return {"gpu_util_mean": round(statistics.mean(s[0] for s in w), 1),
"gpu_util_max": max(s[0] for s in w),
"vram_mb_max": max(s[1] for s in w)}
def one_request(base_url, model, key, max_tokens, thinking):
body = {"model": model,
"messages": [{"role": "system", "content": SYSTEM},
{"role": "user", "content": USER}],
"max_tokens": max_tokens, "temperature": 0.7, "stream": True,
"stream_options": {"include_usage": True},
"chat_template_kwargs": {"enable_thinking": thinking}}
t0 = time.perf_counter(); ttft = None; ctoks = 0
try:
with requests.post(base_url.rstrip("/") + "/chat/completions",
headers={"Authorization": f"Bearer {key}"},
json=body, stream=True, timeout=900) as r:
r.raise_for_status()
for line in r.iter_lines():
if not line: continue
s = line.decode("utf-8", "ignore")
if not s.startswith("data: "): continue
s = s[6:]
if s.strip() == "[DONE]": break
try: ch = json.loads(s)
except json.JSONDecodeError: continue
if ttft is None and ch.get("choices"):
d = ch["choices"][0].get("delta", {})
if d.get("content") or d.get("reasoning_content"):
ttft = time.perf_counter() - t0
if ch.get("usage"):
ctoks = ch["usage"].get("completion_tokens", 0)
return ttft, time.perf_counter() - t0, ctoks, True, None
except Exception as e:
return ttft, time.perf_counter() - t0, ctoks, False, f"{type(e).__name__}: {e}"
def main():
ap = argparse.ArgumentParser()
ap.add_argument("--base-url", required=True)
ap.add_argument("--model", required=True)
ap.add_argument("--api-key", default="sk-harvest-local")
ap.add_argument("--concurrencies", default="1,8,32,64")
ap.add_argument("--rounds", type=int, default=2)
ap.add_argument("--max-tokens", type=int, default=256)
ap.add_argument("--thinking", choices=["on", "off"], default="off")
ap.add_argument("--out", default=None)
a = ap.parse_args()
thinking = a.thinking == "on"
sampler = GpuSampler(); sampler.start()
print(f"model={a.model} thinking={a.thinking} max_tokens={a.max_tokens} rounds={a.rounds}")
print(f"{'conc':>5} {'ok':>7} {'agg tok/s':>10} {'lat p50':>8} {'ttft p50':>9} "
f"{'util avg':>9} {'util max':>9} {'vram MB':>8}")
results = []
for c in [int(x) for x in a.concurrencies.split(",") if x.strip()]:
idx = len(sampler.samples)
n = c * a.rounds
lat, ttfts, toks, oks = [], [], [], 0
w0 = time.perf_counter()
with ThreadPoolExecutor(max_workers=c) as ex:
futs = [ex.submit(one_request, a.base_url, a.model, a.api_key,
a.max_tokens, thinking) for _ in range(n)]
for f in as_completed(futs):
ttft, total, ct, ok, err = f.result()
if ok:
oks += 1; lat.append(total); toks.append(ct)
if ttft is not None: ttfts.append(ttft)
elif err:
print(f" [err] {err[:90]}", file=sys.stderr)
wall = time.perf_counter() - w0
gpu = sampler.window(idx)
agg = sum(toks) / wall if wall else 0
row = {"concurrency": c, "ok": oks, "n": n, "agg_tok_s": round(agg, 1),
"lat_p50_s": round(statistics.median(lat), 2) if lat else None,
"ttft_p50_s": round(statistics.median(ttfts), 3) if ttfts else None,
"mean_completion_tokens": round(statistics.mean(toks), 1) if toks else 0,
"wall_s": round(wall, 1), **gpu}
results.append(row)
print(f"{c:>5} {oks:>3}/{n:<3} {row['agg_tok_s']:>10} {row['lat_p50_s'] or 0:>8} "
f"{row['ttft_p50_s'] or 0:>9} {row.get('gpu_util_mean',0):>9} "
f"{row.get('gpu_util_max',0):>9} {row.get('vram_mb_max',0):>8}")
sampler.on = False
peak = max((r["agg_tok_s"] for r in results), default=0)
print(f"\npeak aggregate: {peak:.0f} tok/s (thinking={a.thinking})")
if a.out:
json.dump({"model": a.model, "thinking": a.thinking,
"max_tokens": a.max_tokens, "results": results,
"peak_agg_tok_s": peak}, open(a.out, "w"), indent=2)
print("wrote", a.out)
if __name__ == "__main__":
main()
|