Skip to content

perf(frontend): merge_chunk rebuilt the whole backlog on every merge - #2020

Merged
valarLip merged 3 commits into
ROCm:mainfrom
JiaoliangYu:jiaolyu/merge-chunk-amortized
Aug 26, 2026
Merged

valarLip merged 3 commits into
ROCm:mainfrom
JiaoliangYu:jiaolyu/merge-chunk-amortized

Conversation

@JiaoliangYu

@JiaoliangYu JiaoliangYu commented Aug 25, 2026 •

Copy link
Copy Markdown
Contributor

Motivation

StreamOutputCollector folds an arriving chunk into the one still waiting whenever the consumer lags. merge_chunk rebuilt token_ids and re-concatenated text from scratch on every fold, which walks the whole accumulation per merge — and it runs on the event loop, so it feeds back on itself: the further the frontend falls behind the engine, the deeper the accumulation, the more each merge costs.

py-spy over 60s of steady decode at c=8192, self time: merge_chunk 27.5% across three lines, the token_ids rebuild alone 20.3% and the largest single consumer in the process.

Technical Details

Both deltas are extended in place. The producer's list is copied once in put_nowait, where the chunk is first taken, so merge_chunk is a single unconditional extend with no per-merge branch.

The copy is defence in depth, not necessity: scheduler.py is the only RequestOutput site in the repo and already hands out a per-step copy, and the object is unpickled fresh after crossing ZMQ. The collector should still not extend a list handed to it across two files and a process boundary on the strength of a comment, and test_put_nowait_gives_merge_chunk_a_list_it_may_extend pins the precondition rather than leaving it to be rediscovered.

text is popped before the concatenation so the string has a single reference and CPython grows it in place; the write-back is in a finally, so a non-string delta raises as it always did without losing the text already accumulated.

Behaviour differences from the previous code, both only reachable by calling merge_chunk outside the collector: token_ids is de-aliased on the way in rather than on the first merge carrying ids, and a chunk that never went through put_nowait raises KeyError instead of silently growing the key.

Test Plan

  • tests/entrypoints/.
  • test_merge_extends_in_place_instead_of_rebuilding and test_put_nowait_gives_merge_chunk_a_list_it_may_extend both fail with the source change reverted.
  • test_a_merged_stream_reads_the_same_as_an_unmerged_one at backlog depths 1, 2, 17 and 500.
  • test_merge_keeps_the_text_it_had_when_the_delta_is_not_a_string.
  • gsm8k 3-shot through lm_eval.
  • benchmark_serving, DSR1 fp4, tp8 with DP attention, c=8192, ISL 1024 / OSL 4096, --ignore-eos, fresh server per run, otherwise identical builds.

lm_eval drives /v1/completions without stream, so gsm8k does not reach the code this PR touches. The backlog-depth cases were written for that.

Test Result

tests/entrypoints/ 1867 passed, 52 skipped, 3 xfailed
source reverted, tests kept 2 failed — the two that pin the change
gsm8k 3-shot 0.9507 flexible / 0.9492 strict, 1319/1319

Throughput, measured on the first implementation of this change (the ownership-tag one); being re-measured on the reworked patch and this section will be replaced with those numbers.

One wave, 8192 requests, --max-num-batched-tokens 4096:

before after
output throughput 42,970 57,735 tok/s +34.4%
mean TPOT 179.86 130.94 ms −27.2%
median ITL 17,475 912 ms −94.8%

Two waves, 16384 requests, default --max-num-batched-tokens:

before after
output throughput 40,873 57,980 tok/s +41.8%
mean TPOT 183.35 130.41 ms −28.9%
median ITL 20,497 1,797 ms −91.2%

Five waves, 40960 requests, 48 minutes: 57,514 tok/s, mean TPOT 131.19 ms, p99 TPOT 140.73 ms, 40959/40960 successful, no traceback and no preemption. One, two and five waves land within 0.8% of each other; before the change more waves meant worse, since the wave boundary is where the accumulation is deepest.

The engine's last decode step lands at 568.6s in both one-wave runs, counted from the decode[bs=...] annotations. The run used to end at 793s and now ends at 581s.

Reverting the engine-side GC tuning from #1980, with merge_chunk untouched, also gains 10.2% on this config: a slower engine keeps the accumulation shallow enough that the per-merge cost stays cheap. Included as evidence for the mechanism, not as an alternative.

Submission Checklist

`StreamOutputCollector` folds an arriving chunk into the one still waiting
whenever the consumer lags, and `merge_chunk` rebuilt `token_ids` and
re-concatenated `text` from scratch each time. That is quadratic in the
backlog depth, and it runs on the event loop: the further the frontend falls
behind the engine the deeper the backlog gets, the more every merge costs,
and the further behind it falls. Past the point where the engine outruns
delivery the frontend does not degrade, it collapses.

The first chunk's list is the engine's own `output_tokens` and still must not
be appended to, so the copy happens once, into a list this module owns and
may extend in place from then on. `text` is popped before the concatenation
so the string has a single reference and CPython resizes it in place instead
of allocating a new one per merge.

DSR1 fp4, tp8 with DP attention, c=8192, ISL 1024 / OSL 4096, one wave,
otherwise identical builds:

                        before      after
  output throughput     42,970  ->  57,735 tok/s   (+34.4%)
  total throughput      53,713  ->  72,169 tok/s
  mean TPOT              179.9  ->   130.9 ms
  median ITL            17,475  ->     912 ms

The engine's last decode step lands at 568.6s in both; the run used to end at
793s and now ends at 581s, so what the quadratic term was buying was 225
seconds of the frontend draining a backlog after the GPUs had gone idle.

Same numbers seen from the other side: reverting the engine-side GC tuning
from ROCm#1980 also gains 10.2% on this config, because a slower engine keeps the
backlog shallow enough that the quadratic term stays cheap. That is worth
noting only as evidence for the mechanism -- the fix belongs here, not there.

gsm8k 3-shot 0.9507 flexible / 0.9492 strict, 1319/1319. Merge output is
byte-identical to the unmerged path at backlog depths from 1 to 500, and to
the previous implementation over 4000 randomized merge sequences, with the
engine's list never mutated.
@github-actions

Copy link
Copy Markdown
Contributor

🏷️ CI Guide

Runs automatically on every eligible PR before approval:

  • ✅ Pre Checkin: Black, Ruff, catalog schema validation, non-GPU unit tests

Heavy model tests:

  • ✅ Run after the PR is approved and Pre Checkin passes
  • ✅ Run immediately when an approval review is submitted
  • ✅ Can be requested before approval with labels
Label Tests
ci:full Run all heavy PR model tests: native ATOM, vLLM, and SGLang
ci:atom Run native ATOM model accuracy tests
ci:vllm Run ATOM vLLM OOT model accuracy tests
ci:sglang Run ATOM SGLang model accuracy tests

Heavy jobs are skipped when the PR is not approved and no matching ci:* label is present.
Add labels via the sidebar or gh pr edit 2020 --add-label <label>

@valarLip

Copy link
Copy Markdown
Collaborator

Nice catch on the quadratic term — the mechanism is real and the depth-500 numbers reproduce for me (633 → 217 ns/merge). Four things before this lands.

1. The premise the copy machinery defends doesn't hold.

The docstring and commit message say the first chunk's list is the engine's own output_tokens and must not be appended to. scheduler.py:2400 is the only RequestOutput( construction in the repo, and lines 2395-2399 already hand out a private per-step copy:

output_tokens_list = list(new_tokens) if isinstance(new_tokens, tuple) else new_tokens.copy()

Sequence.output_tokens is an array("i") (sequence.py:202), and the object is unpickled fresh per step after crossing ZMQ. So _MergedIds, the exact-class check, the conditional copy and the if ids: early-out are all guarding a list nothing else reads. If you want to keep the copy as defence-in-depth that's reasonable, but it should say so, and the coupling should be asserted rather than left implicit 2,400 lines apart.

2. Copying once in put_nowait is simpler and faster at every depth.

merge_chunk has exactly one caller — put_nowait (line 136) — which is also where the producer's list is first stored (self._pending[tag] = chunk, line 134). Adding chunk["token_ids"] = list(chunk.get("token_ids") or ()) there collapses merge_chunk to a single extend:

depth old this PR copy-in-put_nowait
1 297 435 (+46%) 321
2 250 355 (+42%) 259
10 226 254 206
100 341 225 190
500 633 217 186

ns/merge. The shallow-backlog regression matters: depth 1-2 is the common case whenever a merge happens at all, and break-even is around depth 20-30. The variant is already ahead at depth 1 and beats the old code from depth 4 on. Trade-off is honest — it moves the "never touch the producer's list" contract to put_nowait, so test_streaming_dispatch.py:340 would need to drive through there. (I also measured array("i") for consistency with #1990: ~2x slower here, since output_tokens arrives as a list and every merge would pay per-element unboxing. Not worth it.)

3. The docstring now states the opposite of the code.

Lines 81-83 still read "token_ids is rebuilt rather than extended" while line 88 is cur.extend(ids). That's the only prose describing the ownership contract.

The text refcount trick also needs a comment — as written the three lines look like a pointless round-trip, so a cleanup pass rewrites them to into["text"] = into.get("text","") + new.get("text","") and silently restores the quadratic term with no test failure and no lint signal. tool_parser/stream.py:53-68 documents this exact hazard if you want the wording.

4. Three behaviours changed that were structurally impossible before. All reachable only via paths that don't exist today, but they were free guarantees:

merge_chunk({"text": "a"}, {"token_ids": [], "text": "b"})
# -> {'text': 'ab', 'finished': False}   # token_ids key gone; old code left []

eng = [1, 2]
merge_chunk({"token_ids": eng, "text": "a"}, {"text": "b"})
# -> into["token_ids"] is eng == True    # empty-ids path no longer de-aliases

merge_chunk({..., "text": "acc"}, {"token_ids": [2], "text": None})
# -> TypeError, into == {'token_ids': [1, 2], 'finished': False}   # text key popped, never restored

The third is worth fixing regardless — try: text += ... / finally: into["text"] = text keeps the refcount at 1 during the concat. Also note cur.extend turns an aliased or repeated merge from a bounded 2x error into unbounded doubling; a reviewer hit a 2-minute hang on that accidentally.

Tests. The commit message describes exactly the harness this needs — 4000 randomized merge sequences, depths 1 to 500, engine-list-unmutated assertion — but it didn't land, and the suite is green identically with the PR applied and reverted. The existing test_merge_keeps_end_of_stream_and_never_extends_the_engines_list does a single merge, so it only arms the copy branch, never the in-place extend at depth >= 2.

One question on attribution, not a blocker. Backing out from your table: 42,970 tok/s / 8192 = 5.25 tok/s/stream → ~43k merges/s total. At depth 100 the old path is 235 ns/merge ≈ 1.0% of one event-loop core; even at depth 2000 it's 4.3%. That doesn't obviously account for +34% throughput and −95% median ITL. Your own #1980 cross-check hints the dominant cost may have been allocator/GC pressure rather than the copy loop — and the 225 s of post-idle drain is hard to square with a backlog of ≤1 chunk per stream, which is all StreamOutputCollector can hold. If the residual is flush() scheduling one _deliver per engine step into loop._ready with no in-flight check, then the next concurrency step reproduces the collapse and this merge can't touch it. A py-spy or GC-counter delta on the same run would settle it.

@zufayu
zufayu requested a review from yitingw1 August 26, 2026 00:58
…t hold

Review found the premise false. The docstring this change inherited says the
first chunk's list is the engine's own `output_tokens` and must not be appended
to; `scheduler.py` is the only `RequestOutput` site in the repo and already
hands out a per-step copy five lines above it, and the object is unpickled
fresh after crossing ZMQ besides. The ownership tag, the exact-class test and
the conditional copy were all guarding a list nothing else reads.

The copy is kept as defence in depth -- the collector should not extend a list
handed to it across two files and a process boundary on the strength of a
comment -- but it moves to `put_nowait`, which is where the producer's chunk is
first taken and already the only caller. That is one copy per stall rather than
a branch and a type test per merge, and it is faster than the previous
implementation at every backlog depth rather than from depth 20 on.

Three behaviours the old rebuild gave away for free and the tag version had
dropped:

  - a delta with no text no longer loses the text already accumulated. The
    string is popped to keep its refcount at one, so the write-back is in a
    `finally` and a `TypeError` from the concatenation leaves `into` intact.
  - `token_ids` is de-aliased from the producer's list on the way in now,
    unconditionally, rather than on the first merge that carries ids.
  - a chunk that never went through `put_nowait` now raises `KeyError` in
    `merge_chunk` instead of silently growing a key. The precondition is
    pinned by a test rather than left to be rediscovered.

The docstring said "rebuilt rather than extended" over code that extends; it is
rewritten. The pop carries a comment, since restoring `into["text"] = a + b`
reads as a cleanup and brings the quadratic term back with no test or lint
signal.

Tests: the harness the first commit described now exists. Two of the new cases
fail with the source change reverted, which the suite could not do before.
@JiaoliangYu

JiaoliangYu commented Aug 26, 2026 •

Copy link
Copy Markdown
Contributor Author

1. The premise the copy machinery defends doesn't hold.

The docstring and commit message say the first chunk's list is the engine's own output_tokens and must not be appended to. scheduler.py:2400 is the only RequestOutput( construction in the repo, and lines 2395-2399 already hand out a private per-step copy:

output_tokens_list = list(new_tokens) if isinstance(new_tokens, tuple) else new_tokens.copy()

Sequence.output_tokens is an array("i") (sequence.py:202), and the object is unpickled fresh per step after crossing ZMQ. So _MergedIds, the exact-class check, the conditional copy and the if ids: early-out are all guarding a list nothing else reads. If you want to keep the copy as defence-in-depth that's reasonable, but it should say so, and the coupling should be asserted rather than left implicit 2,400 lines apart.

Agreed. _MergedIds, the exact-class check, the conditional copy and the if ids: early-out are all gone.

2. Copying once in put_nowait is simpler and faster at every depth.

merge_chunk has exactly one caller — put_nowait (line 136) — which is also where the producer's list is first stored (self._pending[tag] = chunk, line 134). Adding chunk["token_ids"] = list(chunk.get("token_ids") or ()) there collapses merge_chunk to a single extend:

depth old this PR copy-in-put_nowait
1 297 435 (+46%) 321
2 250 355 (+42%) 259
10 226 254 206
100 341 225 190
500 633 217 186
ns/merge. The shallow-backlog regression matters: depth 1-2 is the common case whenever a merge happens at all, and break-even is around depth 20-30. The variant is already ahead at depth 1 and beats the old code from depth 4 on. Trade-off is honest — it moves the "never touch the producer's list" contract to put_nowait, so test_streaming_dispatch.py:340 would need to drive through there. (I also measured array("i") for consistency with #1990: ~2x slower here, since output_tokens arrives as a list and every merge would pay per-element unboxing. Not worth it.)

Copying once in put_nowait. Adopted verbatim:

if waiting is None:
chunk["token_ids"] = list(chunk.get("token_ids") or ())
self._pending[tag] = chunk

merge_chunk is now a single extend, so the depth-1/2 regression in your table is gone.

3. The docstring now states the opposite of the code.

Lines 81-83 still read "token_ids is rebuilt rather than extended" while line 88 is cur.extend(ids). That's the only prose describing the ownership contract.

The text refcount trick also needs a comment — as written the three lines look like a pointless round-trip, so a cleanup pass rewrites them to into["text"] = into.get("text","") + new.get("text","") and silently restores the quadratic term with no test failure and no lint signal. tool_parser/stream.py:53-68 documents this exact hazard if you want the wording.

The text refcount trick carries a comment saying it is not a pointless round-trip, worded after tool_parser/stream.py:53-68.

4. Three behaviours changed that were structurally impossible before. All reachable only via paths that don't exist today, but they were free guarantees:

merge_chunk({"text": "a"}, {"token_ids": [], "text": "b"})
# -> {'text': 'ab', 'finished': False}   # token_ids key gone; old code left []

eng = [1, 2]
merge_chunk({"token_ids": eng, "text": "a"}, {"text": "b"})
# -> into["token_ids"] is eng == True    # empty-ids path no longer de-aliases

merge_chunk({..., "text": "acc"}, {"token_ids": [2], "text": None})
# -> TypeError, into == {'token_ids': [1, 2], 'finished': False}   # text key popped, never restored

The third is worth fixing regardless — try: text += ... / finally: into["text"] = text keeps the refcount at 1 during the concat. Also note cur.extend turns an aliased or repeated merge from a bounded 2x error into unbounded doubling; a reviewer hit a 2-minute hang on that accidentally.

The third (text=None raising with the key never restored) is fixed with try: / finally:. The other two are now guaranteed by the put_nowait copy — that is the contract migration your point 2 asks for. Stated plainly: calling merge_chunk directly with an into that has no token_ids key now raises KeyError instead of creating it. The unbounded-doubling hazard on aliased or repeated merges is covered by the same copy; into's list is always ours.

Tests. The commit message describes exactly the harness this needs — 4000 randomized merge sequences, depths 1 to 500, engine-list-unmutated assertion — but it didn't land, and the suite is green identically with the PR applied and reverted. The existing test_merge_keeps_end_of_stream_and_never_extends_the_engines_list does a single merge, so it only arms the copy branch, never the in-place extend at depth >= 2.

Five added; two fail against unpatched source, so "green either way" no longer holds. The 4000-sequence harness did not land and I'd rather not add it: test_merging_is_identical_to_unmerged_delivery already sweeps drain cadences 1/2/3/5/2N over a multibyte payload, the equivalence test is parametrized over depths 1/2/17/500, and all four folded fields (extend, +=, or, last-non-empty-wins) are associative.

One question on attribution, not a blocker. Backing out from your table: 42,970 tok/s / 8192 = 5.25 tok/s/stream → ~43k merges/s total. At depth 100 the old path is 235 ns/merge ≈ 1.0% of one event-loop core; even at depth 2000 it's 4.3%. That doesn't obviously account for +34% throughput and −95% median ITL. Your own #1980 cross-check hints the dominant cost may have been allocator/GC pressure rather than the copy loop — and the 225 s of post-idle drain is hard to square with a backlog of ≤1 chunk per stream, which is all StreamOutputCollector can hold. If the residual is flush() scheduling one _deliver per engine step into loop._ready with no in-flight check, then the next concurrency step reproduces the collapse and this merge can't touch it. A py-spy or GC-counter delta on the same run would settle it.

Six arms, same commit, two binary factors (merge implementation, ATOM_GC_THRESHOLD). conc 8192, 8192 requests, 1024/4096, default max-num-batched-tokens. py-spy window pinned at bench+260 s for 60 s; instrumentation byte-identical between arms.

image

Your arithmetic is right, and so is the objection — but merge_chunk's share is not stationary, so no single number for it can be plugged in. py-spy across the whole run, 60 s windows, same commit, only the merge implementation differing:

py-spy self time, merge_chunk bench+60 +180 +300 +420 +540 +660 run-weighted
unpatched 5.4% 7.8% 10.0% 14.8% 27.1% 79.6% 26.9%
patched 3.2% 3.9% 4.3% 4.8% 6.1% (finished) 4.4%

of the one GIL-bound core, the frontend spends 26.9% folding and 73.1% on real work (unpickle, detokenize, deliver) unpatched, against 4.4% and 95.6% patched. The real work per token is unchanged, so capacity rises by 0.956/0.731 = +30.8%, against +32.2% measured on that pair. The model accounts for 96% of the gain; the remainder is plausibly the drain shortening faster than capacity rises.

queueing delay, mean / max unpatched patched
GC set 194.9 / 328.3 s 44.8 / 67.1 s
default 321.8 / 514.2 s 148.7 / 283.6 s

Run duration shrinks by 70–80% of the max-lag reduction in both configurations.

  • GC pressure is not the mechanism. At the default thresholds (gen0 at 97.7/s, 23x the tuned arm without this patch) gen0-per-token moves −0.5% while throughput rises +18.0%.

  • "This merge can't touch it": it does — −77% / −54%. Next concurrency step untested. This doesn't remove the imbalance — it removes the positive feedback. The old path's cost grew with the backlog, so falling behind made every delivery more expensive, which made it fall further behind; that is the 5% → 80% curve. With a constant per-merge cost the backlog grows linearly with the deficit instead of compounding, so the next concurrency step should degrade rather than collapse. The real fixes are backpressure in flush() and cutting the per-token frontend cost — the residual is pickle.loads (~30%) and detokenization (~30%), neither of which this touches — and I'd rather land those separately than widen this PR.

@valarLip
valarLip merged commit bc3ef6e into ROCm:main Aug 26, 2026
34 of 36 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants