From d30b3efa844c6bd6fab94947d915e190d99f5579 Mon Sep 17 00:00:00 2001 From: Bi Xue Date: Fri, 10 Apr 2026 20:52:49 -0700 Subject: [PATCH] [sgl] _ATTN_TP and _ATTN_CP use message queue for broadcast on CPU (#22205) --- python/sglang/srt/distributed/parallel_state.py | 11 ++++------- 1 file changed, 4 insertions(+), 7 deletions(-) diff --git a/python/sglang/srt/distributed/parallel_state.py b/python/sglang/srt/distributed/parallel_state.py index 92e52b10e..98be02687 100644 --- a/python/sglang/srt/distributed/parallel_state.py +++ b/python/sglang/srt/distributed/parallel_state.py @@ -45,7 +45,6 @@ from sglang.srt.compilation.piecewise_context_manager import is_in_piecewise_cud from sglang.srt.distributed.utils import set_global_tcp_store from sglang.srt.environ import envs from sglang.srt.utils import ( - get_bool_env_var, get_current_device_stream_fast, get_int_env_var, is_cpu, @@ -1801,9 +1800,7 @@ def initialize_model_parallel( group_ranks, get_world_group().local_rank, backend, - use_message_queue_broadcaster=get_bool_env_var( - "SGLANG_USE_MESSAGE_QUEUE_BROADCASTER", "true" - ), + use_message_queue_broadcaster=envs.SGLANG_USE_MESSAGE_QUEUE_BROADCASTER.get(), group_name="tp", ) @@ -1816,9 +1813,7 @@ def initialize_model_parallel( group_ranks, get_world_group().local_rank, backend, - use_message_queue_broadcaster=get_bool_env_var( - "SGLANG_USE_MESSAGE_QUEUE_BROADCASTER", "true" - ), + use_message_queue_broadcaster=envs.SGLANG_USE_MESSAGE_QUEUE_BROADCASTER.get(), group_name="pdmux_prefill_tp", ) if _TP.pynccl_comm: @@ -1856,6 +1851,7 @@ def initialize_model_parallel( group_ranks, get_world_group().local_rank, backend, + use_message_queue_broadcaster=envs.SGLANG_USE_MESSAGE_QUEUE_BROADCASTER.get(), group_name="attn_cp", ) @@ -1889,6 +1885,7 @@ def initialize_model_parallel( use_mscclpp_allreduce=False, use_custom_allreduce=False, use_torch_symm_mem_allreduce=False, + use_message_queue_broadcaster=envs.SGLANG_USE_MESSAGE_QUEUE_BROADCASTER.get(), group_name="attention_tp", )