File size: 12,901 Bytes
dbe924d
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
"""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.")