Repository navigation
fix(grpc): forward scheduler load info for DP-aware load balancing - #1115
Kangyan-Zhou wants to merge 1 commit into
Conversation
In gRPC mode, the GrpcRequestManager receives BatchTokenIDOutput from the scheduler but does not forward the piggybacked load info (num_reqs, num_tokens per DP rank) back to the DataParallelController. This causes total_tokens and total_requests load balance policies to see all-zero budgets and degenerate to always picking rank 0. Mirror the TokenizerManager pattern: when dp_size > 1, send a WatchLoadUpdateReq with the load info back to the scheduler input socket so the DataParallelController can make informed dispatch decisions. Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
📝 WalkthroughWalkthroughLoad update forwarding has been added to the batch output handler. When data parallelism is enabled ( Changes
Estimated code review effort🎯 2 (Simple) | ⏱️ ~12 minutes Poem
🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 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 @Kangyan-Zhou, 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 GrpcRequestManager to forward load information to the DataParallelController when data parallelism is enabled, ensuring token-aware balancing functions correctly. The changes involve importing WatchLoadUpdateReq and adding logic to send load updates during batch output processing. A review comment suggests using the internal _send_to_scheduler helper method and awaiting the call to maintain consistency with the rest of the class and ensure proper error logging.
| if self.server_args.dp_size > 1 and batch_out.load is not None: | ||
| load_update = WatchLoadUpdateReq(loads=[batch_out.load]) | ||
| self.send_to_scheduler.send_pyobj(load_update) |
There was a problem hiding this comment.
To maintain consistency with the rest of the GrpcRequestManager class (see lines 368, 427, and 471), please use the _send_to_scheduler helper method instead of calling self.send_to_scheduler.send_pyobj directly. This ensures that the send operation is covered by the helper's error logging. Since _handle_batch_output is an async function, the call should be awaited to follow the established pattern in this class.
| if self.server_args.dp_size > 1 and batch_out.load is not None: | |
| load_update = WatchLoadUpdateReq(loads=[batch_out.load]) | |
| self.send_to_scheduler.send_pyobj(load_update) | |
| if self.server_args.dp_size > 1 and batch_out.load is not None: | |
| load_update = WatchLoadUpdateReq(loads=[batch_out.load]) | |
| await self._send_to_scheduler(load_update) |
References
- If a code block is identified as duplicated across multiple functions or modules, consider refactoring to unify the logic into a shared helper function or class. This improves maintainability and reduces the chance of inconsistencies when changes are needed.
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@grpc_servicer/smg_grpc_servicer/sglang/request_manager.py`:
- Around line 688-690: The send of DP load updates (WatchLoadUpdateReq with
batch_out.load) from _handle_batch_output can raise and abort before put_tasks
are flushed; make this best-effort by catching exceptions around
send_to_scheduler.send_pyobj so failures won't propagate. Wrap the call that
constructs WatchLoadUpdateReq and calls self.send_to_scheduler.send_pyobj(...)
in a try/except that logs/debugs the error (non-fatal) and continues, ensuring
the subsequent flush of put_tasks still runs; keep the guard on
self.server_args.dp_size > 1 and only apply the try/except when batch_out.load
is not None.
🪄 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: f2ce1e0b-ff0e-4497-9371-aab3b5f9d496
📒 Files selected for processing (1)
grpc_servicer/smg_grpc_servicer/sglang/request_manager.py
| if self.server_args.dp_size > 1 and batch_out.load is not None: | ||
| load_update = WatchLoadUpdateReq(loads=[batch_out.load]) | ||
| self.send_to_scheduler.send_pyobj(load_update) |
There was a problem hiding this comment.
Make DP load forwarding best-effort to avoid dropping batch outputs.
At Line 690, an unguarded send can raise and abort _handle_batch_output before Line 694 flushes put_tasks, which can stall/delay client-visible outputs for that batch. Handle this path as non-fatal.
Suggested fix
- if self.server_args.dp_size > 1 and batch_out.load is not None:
- load_update = WatchLoadUpdateReq(loads=[batch_out.load])
- self.send_to_scheduler.send_pyobj(load_update)
+ if self.server_args.dp_size > 1 and batch_out.load is not None:
+ try:
+ await self._send_to_scheduler(WatchLoadUpdateReq(loads=[batch_out.load]))
+ except Exception as e:
+ logger.warning(f"Failed to forward DP load update: {e}")🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@grpc_servicer/smg_grpc_servicer/sglang/request_manager.py` around lines 688 -
690, The send of DP load updates (WatchLoadUpdateReq with batch_out.load) from
_handle_batch_output can raise and abort before put_tasks are flushed; make this
best-effort by catching exceptions around send_to_scheduler.send_pyobj so
failures won't propagate. Wrap the call that constructs WatchLoadUpdateReq and
calls self.send_to_scheduler.send_pyobj(...) in a try/except that logs/debugs
the error (non-fatal) and continues, ensuring the subsequent flush of put_tasks
still runs; keep the guard on self.server_args.dp_size > 1 and only apply the
try/except when batch_out.load is not None.
|
Hi @Kangyan-Zhou, 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:
|
|
Closing in favor of a new PR with correct branch naming convention and DCO signoff. |
Summary
GrpcRequestManagerreceivesBatchTokenIDOutputfrom the scheduler but drops the piggybacked load info (num_reqs,num_tokensper DP rank), causingtotal_tokensandtotal_requestsDP load balance policies to see all-zero budgets and always pick rank 0TokenizerManagerpattern: forwardWatchLoadUpdateReqback to theDataParallelControllerwhendp_size > 1Context
In PD disaggregated serving with DP=8 decode engines, we observed severe DP rank imbalance (e.g., decoder-1 DP2 processed 93,904 tokens while DP7 processed 7). The
round_robinpolicy distributes requests evenly by count but is blind to KV cache pressure from long-running sequences. Switching tototal_tokensrequires this load feedback loop to function.Test plan
--load-balance-method total_tokenson a DP=8 decode engine in gRPC modeDPBudgetreceives non-zero load updates (check scheduler debug logs)sglang_decode_sum_seq_lensunder mixed-length workloadsround_robin(load forwarding is a no-op whendp_size == 1)🤖 Generated with Claude Code
Summary by CodeRabbit