k3-a40-bootstrap / banc_optimisations.py
patdev's picture
balayage optimisations vLLM : prefill, cache, decodage
dbe924d verified
Raw
History Blame Contribute Delete
12.9 kB
"""Balayage des optimisations vLLM : prefill, cache de prefixe, decodage.
Ce que la campagne a etabli, et qui dicte le choix des leviers testes ici :
le plafond solo n'est ni la bande passante (rendement 69 % sur Ada, 22,6 % sur
H200) ni le noyau MoE (`humming` = Marlin a 1 % pres sur trois architectures).
Le terme dominant est un COUT FIXE PAR JETON. Les leviers qui peuvent le
reduire sont donc : la compilation, la couverture des graphes CUDA,
l'ordonnancement, et le noyau LINEAIRE -- que `--moe-backend` ne touche pas.
Detail decisif releve dans les journaux precedents : meme avec
`--moe-backend humming`, vLLM affiche toujours
`Using MarlinNvFp4LinearKernel for NVFP4 GEMM`. Le drapeau ne remplace que les
experts. Or le chemin dense pese 1,849 Gio des 2,880 Gio du socle actif, soit
64 %. On n'avait donc echange que 36 % du travail.
Trois metriques par configuration, et non une :
- PREFILL : jetons de prompt / TTFT sur un prompt froid, abscisse lue dans
usage.prompt_tokens et jamais estimee ;
- CACHE : taux de reussite du cache de prefixe, lu dans /metrics, plus le
TTFT du meme prompt rejoue ;
- DECODAGE : debit par flux a 1 session et agrege a 8.
"""
import json
import os
import re
import statistics
import subprocess
import sys
import threading
import time
import urllib.request
MODEL = "nvidia/NVIDIA-Nemotron-3.5-Lightning-30B-A3B-NVFP4"
PORT = 8000
URL = "http://127.0.0.1:%d" % PORT
def dire(*a):
print(*a, flush=True)
def titre(t):
dire("\n" + "=" * 74)
dire(t)
dire("=" * 74)
titre("0. la carte, la pile, et la selection du noyau LINEAIRE")
subprocess.run(["nvidia-smi", "--query-gpu=name,memory.total,compute_cap",
"--format=csv,noheader"], check=False)
subprocess.run(["python3", "-c",
"import vllm,torch;print('vllm',vllm.__version__,'torch',torch.__version__,"
"'cap',torch.cuda.get_device_capability(0))"], check=False)
# Sonde gratuite : existe-t-il un moyen de changer le noyau NVFP4 LINEAIRE ?
# `--moe-backend` ne le touche pas, et c'est lui qui porte 64 % du socle.
dire("\n--- noyaux NVFP4 lineaires enregistres dans vLLM ---")
sonde = r'''
import inspect, os, re
try:
from vllm.model_executor.layers.quantization.kernels.mixed_precision import __init__ as _
except Exception:
pass
trouve = []
import vllm, pathlib
racine = pathlib.Path(vllm.__file__).parent
for p in racine.rglob("*.py"):
try:
t = p.read_text(errors="ignore")
except Exception:
continue
if "NvFp4LinearKernel" in t or "NVFP4 GEMM" in t:
for m in re.finditer(r"class\s+(\w*NvFp4\w*Kernel)\b", t):
trouve.append((m.group(1), str(p.relative_to(racine))))
for m in re.finditer(r'VLLM_[A-Z0-9_]*(?:NVFP4|GEMM|LINEAR)[A-Z0-9_]*', t):
trouve.append(("env:" + m.group(0), str(p.relative_to(racine))))
vus = set()
for nom, ou in trouve:
if nom in vus:
continue
vus.add(nom)
print(" %-42s %s" % (nom, ou))
if not vus:
print(" aucun -- la selection est probablement en dur")
'''
subprocess.run(["python3", "-c", sonde], check=False)
# Recette NVIDIA comme socle commun ; chaque essai n'ajoute que sa variante.
BASE = ["vllm", "serve", MODEL,
"--served-model-name", "ornith",
"--host", "127.0.0.1", "--port", str(PORT),
"--trust-remote-code",
"--max-model-len", "131072",
"--moe-backend", "marlin",
"--kv-cache-dtype", "fp8",
"--enable-prefix-caching",
"--gpu-memory-utilization", "0.85",
"--mamba-backend", "flashinfer",
"--mamba-cache-mode", "align",
"--reasoning-parser", "nemotron_v3",
"--tool-call-parser", "qwen3_coder",
"--enable-auto-tool-choice"]
# Un levier par ligne, pour qu'un gain soit attribuable. Le dernier essai
# combine les gagnants -- et c'est le seul dont le resultat n'est pas
# attribuable a un levier unique.
ESSAIS = [
("reference (recette NVIDIA)", [], {}),
("ordonnancement asynchrone", ["--async-scheduling"], {}),
("compilation -O3", ["-O3"], {}),
("graphes CUDA FULL", ["--compilation-config", '{"cudagraph_mode":"FULL"}'], {}),
("attention TRITON_ATTN", ["--attention-backend", "TRITON_ATTN"], {}),
("mamba-cache-mode all", ["--mamba-cache-mode", "all"], {}),
("KV en auto (pas fp8)", ["--kv-cache-dtype", "auto"], {}),
("max-num-seqs 8 (mono-session)", ["--max-num-seqs", "8"], {}),
("max-num-batched-tokens 8192", ["--max-num-batched-tokens", "8192"], {}),
("sans cache de prefixe", ["--no-enable-prefix-caching"], {}),
("echantillonneur flashinfer coupe", [], {"VLLM_USE_FLASHINFER_SAMPLER": "0"}),
]
LONG = ("Voici un module Python a auditer.\n\n" +
"\n".join("def f%d(x):\n # etape %d du pipeline de traitement\n"
" y = x * %d + %d\n return y if y > 0 else -y\n" % (i, i, i % 7 + 1, i)
for i in range(1400)))
def demarrer(sup, env_sup, journal):
env = dict(os.environ)
env.update(env_sup)
with open(journal, "w") as f:
p = subprocess.Popen(BASE + sup, stdout=f, stderr=subprocess.STDOUT, env=env)
for i in range(80):
try:
urllib.request.urlopen(URL + "/v1/models", timeout=5).read()
return p, i * 10
except Exception:
pass
if p.poll() is not None:
return None, i * 10
time.sleep(10)
p.terminate()
return None, 800
def appel(contenu, sortie=300, stream=True):
corps = json.dumps({"model": "ornith",
"messages": [{"role": "user", "content": contenu}],
"max_tokens": sortie, "temperature": 0.0,
"stream": stream,
"stream_options": {"include_usage": True} if stream else None
}).encode()
r = urllib.request.Request(URL + "/v1/chat/completions", data=corps,
headers={"Content-Type": "application/json"})
t0 = time.time()
t1 = None
n = 0
bouts = []
prompt_tokens = None
with urllib.request.urlopen(r, timeout=900) as rep:
for l in rep:
l = l.strip()
if not l.startswith(b"data: ") or l[6:] == b"[DONE]":
continue
d = json.loads(l[6:])
if d.get("usage"):
prompt_tokens = d["usage"].get("prompt_tokens")
ch = (d.get("choices") or [{}])
if not ch:
continue
de = ch[0].get("delta", {}) or {}
x = de.get("content") or de.get("reasoning") or de.get("reasoning_content")
if x:
if t1 is None:
t1 = time.time()
n += 1
bouts.append(x)
return {"ttft": (t1 - t0) if t1 else None, "n": n,
"t1": t1, "t2": time.time(), "prompt_tokens": prompt_tokens,
"txt": "".join(bouts)}
def metriques():
try:
t = urllib.request.urlopen(URL + "/metrics", timeout=15).read().decode()
except Exception:
return {}
out = {}
for cle in ("prefix_cache_queries_total", "prefix_cache_hits_total"):
m = re.search(r"vllm:gpu_%s\{[^}]*\}\s+([0-9.e+]+)" % cle, t) or \
re.search(r"vllm:%s\{[^}]*\}\s+([0-9.e+]+)" % cle, t)
if m:
out[cle] = float(m.group(1))
return out
def div4(t):
m = t.split()
if len(m) < 40:
return 1.0
g = [tuple(m[i:i + 4]) for i in range(len(m) - 3)]
return len(set(g)) / len(g)
SUJETS = ["un cache LRU avec dict et liste doublement chainee",
"un pool de connexions avec expiration et sante des sockets",
"un analyseur d'expressions par descente recursive",
"une file de priorite par tas binaire",
"un limiteur de debit par seau a jetons",
"un index inverse pour recherche plein texte",
"un ordonnanceur de taches avec dependances",
"un serialiseur binaire versionne"]
def decodage(conc):
res = [None] * conc
def un(i):
try:
res[i] = appel("Ecris en Python %s, avec trois tests unittest."
% SUJETS[i % len(SUJETS)])
except Exception as e:
res[i] = {"err": str(e)[:60]}
d0 = time.time()
fils = [threading.Thread(target=un, args=(i,)) for i in range(conc)]
for f in fils:
f.start()
for f in fils:
f.join()
d1 = time.time()
bons = [r for r in res if r and not r.get("err") and r.get("t1")]
if not bons:
return None
return {"agrege": sum(r["n"] for r in bons) / (d1 - d0),
"par_flux": statistics.median([(r["n"] - 1) / (r["t2"] - r["t1"])
for r in bons if r["t2"] > r["t1"]]),
"div4": statistics.median([div4(r["txt"]) for r in bons])}
resume = []
for idx, (etiq, sup, env_sup) in enumerate(ESSAIS):
titre("%d. %s" % (idx + 1, etiq))
journal = "/tmp/opt_%d.log" % idx
proc, secondes = demarrer(sup, env_sup, journal)
texte = open(journal, errors="replace").read()
for motif in ("NvFp4 MoE backend", "NVFP4 GEMM", "GPU KV cache size",
"attention backend", "Capturing", "cudagraph"):
for ligne in texte.splitlines():
if motif in ligne:
dire(" " + ligne.split("] ")[-1][:145])
break
if not proc:
dire(" NE DEMARRE PAS (%d s)" % secondes)
vu = set()
for ligne in texte.splitlines():
if any(m in ligne for m in ("RuntimeError", "ValueError", "Traceback",
"unrecognized arguments", "invalid choice",
"NotImplementedError", "AssertionError")):
t = ligne.split("] ")[-1][:160]
if t not in vu:
vu.add(t)
dire(" > " + t)
if len(vu) >= 4:
break
resume.append((etiq, None, None, None, None))
continue
dire(" PRET en %d s" % secondes)
# --- PREFILL, sur un prompt froid, abscisse LUE ---
m0 = metriques()
try:
froid = appel(LONG, sortie=16)
pt = froid["prompt_tokens"]
pf = (pt / froid["ttft"]) if (pt and froid["ttft"]) else None
dire(" prefill froid : %s jetons en %.2f s -> %s jetons/s"
% (pt, froid["ttft"] or 0, ("%.0f" % pf) if pf else "?"))
except Exception as e:
pt = pf = None
froid = {"ttft": None}
dire(" prefill froid : ECHEC %s" % str(e)[:70])
# --- CACHE DE PREFIXE : le meme prompt rejoue ---
gain = None
try:
chaud = appel(LONG, sortie=16)
if froid.get("ttft") and chaud.get("ttft"):
gain = froid["ttft"] / chaud["ttft"]
dire(" prefill chaud : %.2f s (x%.1f plus rapide)"
% (chaud["ttft"], gain))
except Exception as e:
dire(" prefill chaud : ECHEC %s" % str(e)[:70])
m1 = metriques()
taux = None
if m1.get("prefix_cache_queries_total") and m0 is not None:
dq = m1.get("prefix_cache_queries_total", 0) - m0.get("prefix_cache_queries_total", 0)
dh = m1.get("prefix_cache_hits_total", 0) - m0.get("prefix_cache_hits_total", 0)
if dq > 0:
taux = dh / dq
dire(" cache de prefixe : %.1f %% de reussite (%.0f/%.0f jetons)"
% (taux * 100, dh, dq))
if taux is None:
dire(" cache de prefixe : compteurs absents")
# --- DECODAGE ---
dire("conc | agrege | par flux | 4-gr")
solo = agg8 = None
for conc in (1, 8):
d = decodage(conc)
if not d:
dire("%4d | ECHEC" % conc)
continue
dire("%4d | %8.1f | %8.1f | %.3f %s"
% (conc, d["agrege"], d["par_flux"], d["div4"],
"" if d["div4"] > 0.6 else " DEGENERE"))
if conc == 1:
solo = d["par_flux"]
else:
agg8 = d["agrege"]
resume.append((etiq, solo, agg8, pf, taux))
proc.terminate()
time.sleep(20)
titre("RESUME")
dire("%-34s %8s %9s %10s %8s" % ("configuration", "solo", "agrege@8", "prefill/s", "cache"))
ref = resume[0][1] if resume and resume[0][1] else None
for etiq, solo, agg8, pf, taux in resume:
d = ""
if ref and solo:
d = " %+5.1f %%" % (100 * (solo - ref) / ref)
dire("%-34s %8s %9s %10s %8s%s"
% (etiq,
("%.1f" % solo) if solo else "echec",
("%.0f" % agg8) if agg8 else "-",
("%.0f" % pf) if pf else "-",
("%.0f %%" % (taux * 100)) if taux is not None else "-",
d))
dire("\nUn levier n'est retenu que si son gain depasse la dispersion entre deux")
dire("executions identiques -- environ 2 % sur ce banc. En dessous, c'est du bruit.")