Repository navigation
fix(cache-aware): fix load imbalance when decode is faster than prefill - #1714
Conversation
|
Note Reviews pausedIt looks like this branch is under active development. To avoid overwhelming you with review comments due to an influx of new commits, CodeRabbit has automatically paused this review. You can configure this behavior by changing the Use the following commands to manage reviews:
Use the checkboxes below for quick actions:
No actionable comments were generated in the recent review. 🎉 ℹ️ Recent review info⚙️ Run configurationConfiguration used: Organization UI Review profile: ASSERTIVE Plan: Pro Run ID: 📒 Files selected for processing (2)
📝 WalkthroughWalkthroughWorker selection in the cache-aware policy now breaks ties using processed request counts alongside load. Load guard management in the PD router's streaming response path is refactored to construct guards upfront and pass them explicitly through internal helpers, rather than creating them conditionally inside response constructors. ChangesWorker Selection and Load Guard Flow
Estimated code review effort🎯 3 (Moderate) | ⏱️ ~22 minutes Suggested labels
Suggested reviewers
Poem
🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
|
Hi @SYChen123, the DCO sign-off check has failed. All commits must include a To fix existing commits: # Sign off the last N commits (replace N with the number of unsigned commits)
git rebase HEAD~N --signoff
git push --force-with-leaseTo sign off future commits automatically:
|
There was a problem hiding this comment.
Code Review
This pull request updates the worker selection logic in cache_aware.rs to tie-break using processed requests and worker indices, and refactors load management in pd_router.rs by eagerly creating and passing WorkerLoadGuards. The feedback suggests extracting the duplicated worker tie-breaking closure logic into a shared helper function to improve maintainability, and adding a trailing comma in the test file to match standard Rust formatting.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: f81b8665ed
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@model_gateway/src/policies/cache_aware.rs`:
- Around line 486-487: The event-driven no-overlap fallback still chooses by
workers[idx].load() only; update that fallback to use the same tie-break key as
the other branches — (workers[idx].load(), workers[idx].processed_requests(),
idx) — so selection matches the min_by_key change, and add a regression test
(either test_event_driven_no_overlap_uses_min_load or
test_event_driven_short_request_uses_min_load) that exercises the gRPC+KV events
cold-start/no-overlap path to assert selection prefers lower processed_requests
on equal load.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 0d211090-1cf0-4d12-b659-e6399f37eab2
📒 Files selected for processing (2)
model_gateway/src/policies/cache_aware.rsmodel_gateway/src/routers/http/pd_router.rs
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 1ad3d475b3
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@model_gateway/src/routers/http/pd_router.rs`:
- Around line 1578-1580: The call to handle_decode_error_response has arguments
in the wrong positions and types. The function expects the parameters in this
order: response, context reference, decode worker, and then load guards as a
Vec. The current code incorrectly passes prefill as the third argument instead
of decode, and passes decode (an Arc<dyn Worker>) as the fourth argument instead
of the required Vec<WorkerLoadGuard>. Fix this by passing decode as the third
argument and creating a Vec containing WorkerLoadGuard instances for both the
prefill and decode workers using WorkerLoadGuard::new() for each worker as the
fourth argument.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 06d1a442-209a-4753-a0a5-98352fae0697
📒 Files selected for processing (2)
model_gateway/src/policies/cache_aware.rsmodel_gateway/src/routers/http/pd_router.rs
There was a problem hiding this comment.
Caution
Inline review comments failed to post. This is likely due to GitHub's internal server error or limits when posting large numbers of comments. If you are seeing this consistently it is likely a permissions issue. Please check "Moderation" -> "Code review limits" under your organization settings.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@model_gateway/src/routers/http/pd_router.rs`:
- Around line 1578-1580: The call to handle_decode_error_response has arguments
in the wrong positions and types. The function expects the parameters in this
order: response, context reference, decode worker, and then load guards as a
Vec. The current code incorrectly passes prefill as the third argument instead
of decode, and passes decode (an Arc<dyn Worker>) as the fourth argument instead
of the required Vec<WorkerLoadGuard>. Fix this by passing decode as the third
argument and creating a Vec containing WorkerLoadGuard instances for both the
prefill and decode workers using WorkerLoadGuard::new() for each worker as the
fourth argument.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 06d1a442-209a-4753-a0a5-98352fae0697
📒 Files selected for processing (2)
model_gateway/src/policies/cache_aware.rsmodel_gateway/src/routers/http/pd_router.rs
🛑 Comments failed to post (1)
model_gateway/src/routers/http/pd_router.rs (1)
1578-1580:
⚠️ Potential issue | 🔴 Critical🧩 Analysis chain
🏁 Script executed:
#!/bin/bash # Verify if the test compiles with the current code cd model_gateway && cargo test --no-run --test-threads=1 2>&1 | head -50Repository: lightseekorg/smg
Length of output: 307
🏁 Script executed:
# Read the function signature around line 488-494 sed -n '488,494p' model_gateway/src/routers/http/pd_router.rsRepository: lightseekorg/smg
Length of output: 286
🏁 Script executed:
# Read the test call around line 1578-1580 sed -n '1575,1585p' model_gateway/src/routers/http/pd_router.rsRepository: lightseekorg/smg
Length of output: 409
🏁 Script executed:
# Find the variable declarations for prefill and decode in the test (around line 1550-1560) sed -n '1545,1565p' model_gateway/src/routers/http/pd_router.rsRepository: lightseekorg/smg
Length of output: 724
🏁 Script executed:
# Find WorkerLoadGuard definition to understand the type rg "struct WorkerLoadGuard|enum WorkerLoadGuard" -A 5 model_gateway/src/Repository: lightseekorg/smg
Length of output: 387
🏁 Script executed:
# Check if there are any imports or type aliases that might affect the types sed -n '1,50p' model_gateway/src/routers/http/pd_router.rsRepository: lightseekorg/smg
Length of output: 1535
Type mismatch: test passes
Arc<dyn Worker>whereVec<WorkerLoadGuard>is expected.The function signature at line 488-494 shows the parameters are
res: reqwest::Response,context: &PDRequestContext<'_>,decode: Arc<dyn Worker>, andload_guards: Vec<WorkerLoadGuard>. However, the test at lines 1578-1580 passesdecode_response,&context,prefill, anddecodeas arguments. The fourth argument is anArc<dyn Worker>but the function expects aVec<WorkerLoadGuard>. Additionally, passingprefillto thedecodeparameter is semantically incorrect—the decode worker should be passed instead.Create load guards using
WorkerLoadGuard::new()for both workers and pass them as the fourth argument, withdecodeas the third argument.🐛 Proposed fix
- let response = router - .handle_decode_error_response(decode_response, &context, prefill, decode) - .await; + let load_guards = vec![ + WorkerLoadGuard::new(prefill.clone(), None), + WorkerLoadGuard::new(decode.clone(), None), + ]; + + let response = router + .handle_decode_error_response(decode_response, &context, decode, load_guards) + .await;📝 Committable suggestion
‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.let load_guards = vec![ WorkerLoadGuard::new(prefill.clone(), None), WorkerLoadGuard::new(decode.clone(), None), ]; let response = router .handle_decode_error_response(decode_response, &context, decode, load_guards) .await;🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@model_gateway/src/routers/http/pd_router.rs` around lines 1578 - 1580, The call to handle_decode_error_response has arguments in the wrong positions and types. The function expects the parameters in this order: response, context reference, decode worker, and then load guards as a Vec. The current code incorrectly passes prefill as the third argument instead of decode, and passes decode (an Arc<dyn Worker>) as the fourth argument instead of the required Vec<WorkerLoadGuard>. Fix this by passing decode as the third argument and creating a Vec containing WorkerLoadGuard instances for both the prefill and decode workers using WorkerLoadGuard::new() for each worker as the fourth argument.
…h slower than decode Signed-off-by: Simo Lin <25425177+slin1237@users.noreply.github.com>
…drop dead decode param Apply the (load, processed_requests, idx) tiebreak to the 4th least-load site (the event-driven KV-aware 'no overlap' fallback at cache_aware.rs) that the original change missed, so all least-load selections spread evenly when decode is faster than prefill. Remove the now-unused 'decode' parameter from create_streaming_response (Change B moved guard creation out but left it, failing clippy -D warnings) and its call sites. Run cargo +nightly fmt --all (fixes the mis-indented min_by_key closures + trailing whitespace). Signed-off-by: Simo Lin <25425177+slin1237@users.noreply.github.com>
Signed-off-by: Simo Lin <25425177+slin1237@users.noreply.github.com>
… test The empty-indexer fallthrough test routed a 4-token input, but PAGE_SIZE is 16, so the sequence was never cacheable: nothing was inserted and both requests took the min-load branch. It only passed before because the load-only tiebreak was stable by index. With the new (load, processed_requests, idx) tiebreak, identical uncacheable requests correctly spread across workers, so the test must use a >= PAGE_SIZE sequence to exercise a genuine cache hit (which routes by tenant, not load). Signed-off-by: Simo Lin <25425177+slin1237@users.noreply.github.com>
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 3228427ef3
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| .iter() | ||
| .min_by_key(|&&idx| workers[idx].load()) | ||
| .min_by_key(|&&idx| { | ||
| (workers[idx].load(), workers[idx].processed_requests(), idx) |
There was a problem hiding this comment.
Reserve worker before tree insertion
When several cold cache-miss requests are selected concurrently on different Tokio worker threads, this processed_requests() tie-breaker is not made visible until after match_and_insert_with finishes and increment_processed() runs below, while the WorkerLoadGuard is only created after select_pd_pair returns. For long prompts, multiple requests can therefore observe identical (load, processed_requests) values, all choose the lowest index, and insert for that same worker—the imbalance scenario this change is trying to avoid. Reserve or increment the chosen worker before the tree insertion, and apply the same ordering to the token/min-load paths.
Useful? React with 👍 / 👎.
select_worker made several O(workers) passes per request (the healthy filter, the is_imbalanced load fold, and the cache-hit/miss worker scans), and each per-worker access — status, circuit breaker, load — took its own arc_swap guard. At high worker counts that per-worker guard traffic dominated routing CPU. Read each worker once via a new Worker::routing_state() that shares the runtime guard for status+load+processed, gathering the healthy set, load min/max and the min-load index in a single pass. is_imbalanced and the min-load fallback consume the gathered bounds; the cache-hit tenant lookup is a hash-free scan over the gathered healthy indices (url() is a cheap field read). The (load, processed_requests, idx) min-load tie-break from #1714 rides the same guard, so it costs nothing extra. Selection is unchanged except that a cache hit no longer routes onto a Ready-but-circuit-broken worker (it falls through to min-load, like the rest of the selection). No-GPU sim (4 threads, 2048 HTTP workers, shared-prefix load): cache_aware routing cost over round_robin drops ~2x at scale. 26 cache_aware + 107 policy + 207 worker tests pass. Signed-off-by: Simo Lin <25425177+slin1237@users.noreply.github.com>
select_worker made several O(workers) passes per request (the healthy filter, the is_imbalanced load fold, and the cache-hit/miss worker scans), and each per-worker access — status, circuit breaker, load — took its own arc_swap guard. At high worker counts that per-worker guard traffic dominated routing CPU. Read each worker once via a new Worker::routing_state() that shares the runtime guard for status+load+processed, gathering the healthy set, load min/max and the min-load index in a single pass. is_imbalanced and the min-load fallback consume the gathered bounds; the cache-hit tenant lookup is a hash-free scan over the gathered healthy indices (url() is a cheap field read). The (load, processed_requests, idx) min-load tie-break from #1714 rides the same guard, so it costs nothing extra. Selection is unchanged except that a cache hit no longer routes onto a Ready-but-circuit-broken worker (it falls through to min-load, like the rest of selection). No-GPU sim (4 threads, 2048 HTTP workers, shared-prefix load), A/B back-to-back: cache_aware routing is now within noise of round_robin at 2048 workers (+0.024 cpu_ms/req marginal, vs +0.18 with a per-request url map and ~+0.40 unoptimized). 26 cache_aware + 107 policy + 207 worker tests pass. Signed-off-by: Simo Lin <25425177+slin1237@users.noreply.github.com>
select_worker made several O(workers) passes per request (the healthy filter, the is_imbalanced load fold, and the cache-hit/miss worker scans), and each per-worker access — status, circuit breaker, load — took its own arc_swap guard. At high worker counts that per-worker guard traffic dominated routing CPU. Read each worker once via a new Worker::routing_state() that shares the runtime guard for status+load+processed, gathering the healthy set, load min/max and the min-load index in a single pass. is_imbalanced and the min-load fallback consume the gathered bounds; the cache-hit tenant lookup is a hash-free scan over the gathered healthy indices (url() is a cheap field read). The (load, processed_requests, idx) min-load tie-break from #1714 rides the same guard, so it costs nothing extra. Selection is unchanged except that a cache hit no longer routes onto a Ready-but-circuit-broken worker (it falls through to min-load, like the rest of selection). No-GPU sim (4 threads, 2048 HTTP workers, shared-prefix load), A/B back-to-back: cache_aware routing is now within noise of round_robin at 2048 workers (+0.024 cpu_ms/req marginal, vs +0.18 with a per-request url map and ~+0.40 unoptimized). 26 cache_aware + 107 policy + 207 worker tests pass. Signed-off-by: Simo Lin <25425177+slin1237@users.noreply.github.com>
|
@slin1237 Thanks for your review and merge. Could you please merge the PR with the same modification in sglang as well? So that this fix can take effect in next version of sglang. Thanks a lot! |
Description
Problem
In some extreme scenarios, default cache_aware policy will lead to severe load imbalance problem.
For example, 32k input and output only 1 token. In such case, all the requests will be routed to the same prefill worker and others are always left unused.
The reason is that in the select_worker of CacheAwarePolicy, only workers[idx].load() is used for sorting. After the worker is selected, worker.load() will only be updated till create_streaming_response is called (after stream response is returned from decode worker to the router). If prefill is very slow and decode is very fast, then at most of the time, worker.load() will be the same for all prefill workers. Then pd-router will always select the first prefill worker.
Solution
Changes
Test Plan
cd ./model-gateway && TMPDIR=/private/tmp cargo test --lib test_streaming_load_trackingbenching random dataset with 32k input and 1 token output, all the prefill workers have requests received.
Checklist
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspassesSummary by CodeRabbit
Summary of changes
Performance Improvements
Bug Fixes
Refactor
Tests