[metrics 2/5] Report aggregate engine throughput imbalance - #2366
Conversation
c841de6 to
ee74dad
Compare
ee74dad to
1ddc90c
Compare
1ddc90c to
b64172a
Compare
b64172a to
d3bf365
Compare
There was a problem hiding this comment.
Code Review
This pull request introduces cross-engine throughput imbalance tracking to the vLLM metrics scraper by calculating the coefficient of variation (CV) for prompt and generation token deltas. It also adds corresponding unit tests to verify the imbalance calculation under various scenarios. The review feedback suggests two performance optimizations: filtering metrics by name early in the parsing loop to avoid redundant dictionary allocations, and caching the result of statistics.mean(deltas) to prevent duplicate calculations.
| engine_counters = {} | ||
| for (name, labels), value in parsed.items(): | ||
| label_dict = dict(labels) | ||
| worker = label_dict.get("WorkerId", label_dict.get("ReplicaId")) | ||
| if worker is None: | ||
| continue | ||
| engine = (worker, label_dict.get("engine", "0")) | ||
| if name == _GAUGE_NUM_RUNNING: | ||
| engine_counters.setdefault(engine, {}) | ||
| elif name in (_COUNTER_PROMPT_TOKENS, _COUNTER_GENERATION_TOKENS): | ||
| counters = engine_counters.setdefault(engine, {}) | ||
| counters[name] = counters.get(name, 0.0) + value | ||
| self._engine_snapshot = engine_counters |
There was a problem hiding this comment.
For efficiency, we should filter by the metric names we care about (_GAUGE_NUM_RUNNING, _COUNTER_PROMPT_TOKENS, _COUNTER_GENERATION_TOKENS) before converting the labels to a dictionary and extracting the worker/engine information. Since parsed can contain many other metrics (such as numerous histogram buckets), this avoids unnecessary dictionary allocations and lookups in a hot path.
| engine_counters = {} | |
| for (name, labels), value in parsed.items(): | |
| label_dict = dict(labels) | |
| worker = label_dict.get("WorkerId", label_dict.get("ReplicaId")) | |
| if worker is None: | |
| continue | |
| engine = (worker, label_dict.get("engine", "0")) | |
| if name == _GAUGE_NUM_RUNNING: | |
| engine_counters.setdefault(engine, {}) | |
| elif name in (_COUNTER_PROMPT_TOKENS, _COUNTER_GENERATION_TOKENS): | |
| counters = engine_counters.setdefault(engine, {}) | |
| counters[name] = counters.get(name, 0.0) + value | |
| self._engine_snapshot = engine_counters | |
| engine_counters = {} | |
| for (name, labels), value in parsed.items(): | |
| if name not in (_GAUGE_NUM_RUNNING, _COUNTER_PROMPT_TOKENS, _COUNTER_GENERATION_TOKENS): | |
| continue | |
| label_dict = dict(labels) | |
| worker = label_dict.get("WorkerId", label_dict.get("ReplicaId")) | |
| if worker is None: | |
| continue | |
| engine = (worker, label_dict.get("engine", "0")) | |
| if name == _GAUGE_NUM_RUNNING: | |
| engine_counters.setdefault(engine, {}) | |
| else: | |
| counters = engine_counters.setdefault(engine, {}) | |
| counters[name] = counters.get(name, 0.0) + value | |
| self._engine_snapshot = engine_counters |
| if len(deltas) >= 2 and statistics.mean(deltas) > 0: | ||
| result[prefix + public] = statistics.pstdev(deltas) / statistics.mean(deltas) |
There was a problem hiding this comment.
Avoid calculating statistics.mean(deltas) twice by storing the result in a local variable.
| if len(deltas) >= 2 and statistics.mean(deltas) > 0: | |
| result[prefix + public] = statistics.pstdev(deltas) / statistics.mean(deltas) | |
| if len(deltas) >= 2: | |
| mean = statistics.mean(deltas) | |
| if mean > 0: | |
| result[prefix + public] = statistics.pstdev(deltas) / mean |
|
There was a problem hiding this comment.
Cursor Bugbot has reviewed your changes and found 1 potential issue.
❌ Bugbot Autofix is OFF. To automatically fix reported issues with cloud agents, have a team admin enable autofix in the Cursor dashboard.
Reviewed by Cursor Bugbot for commit 915eaa3. Configure here.
915eaa3 to
b0be002
Compare
Signed-off-by: Codex <codex@users.noreply.github.com> Signed-off-by: SumanthRH <sumanthrh99@gmail.com>
Signed-off-by: SumanthRH <sumanthrh99@gmail.com>
b0be002 to
e9e9edf
Compare

What does this PR do?
Report how evenly inference engines share prompt and output generation work.
TLDR: Add aggregate coefficients of variation (CV) for per-engine token deltas, including observed idle engines.
How it works
Group token counters by WorkerId and engine index, then compare each engine's deltas over the sampling interval. CV is population standard deviation divided by mean; a common time denominator cancels from this ratio.
Emit
prompt_throughput_cvandgeneration_throughput_cv, plus a_num_enginesfield for each. Engines observed at both boundaries contribute; invalid/reset deltas are excluded. Omit CV for fewer than two contributing engines or zero mean. Sync windows require a valid start snapshot before deriving CV or engine counts; failed baselines leave only the current gauges. The training logger receives aggregate scalars.Validation
Focused cases check idle-engine inclusion, equal load, missing engines and zero work.
Combined-stack CPU check at
04290d61(#2401): 70 tests passed.--noconftestskips the repository's automatic Ray lifecycle fixture so these mock/CPU checks do not attach to the live cluster. Ray 2.58.0 matches the validation cluster.The existing HTTP503 reproducer confirms that a failed sync baseline returns current gauges without CV or engine counts. All 29 existing collector, histogram and imbalance tests also passed on this branch. No tests were added or edited, and no GPU benchmark was run.
Note
Low Risk
Observability-only additions to metrics aggregation with guarded CV math; no changes to inference or training control flow.
Overview
Adds per-engine load imbalance signals to vLLM Ray metrics scraping by tracking prompt and generation token counters per
(WorkerId/ReplicaId, engine)and emitting aggregate coefficients of variation over each sampling window.VLLMMetricsScrapernow builds an_engine_snapshoton each scrape (including idle engines seen vianum_requests_running), keeps previous/window baselines, and mergesprompt_throughput_cv,generation_throughput_cv, and matching_num_enginesinto both stepsample()output (vllm/…) and explicitstart/stopwindows ({label}/…). CV is computed from non-negative token deltas across engines present at both boundaries; it is omitted when fewer than two engines contribute or mean delta is zero.New unit tests cover idle-engine inclusion, equal load (CV=0), and missing/zero-work cases for
engine_imbalance.Reviewed by Cursor Bugbot for commit e9e9edf. Bugbot is set up for automated code reviews on this repo. Configure here.