Replace all-reduce + dp_scatter with reduce_scatterv for DP attention (#22642)
This commit is contained in:
@@ -47,6 +47,7 @@ from sglang.srt.layers.dp_attention import (
|
|||||||
get_attention_dp_size,
|
get_attention_dp_size,
|
||||||
get_attention_tp_rank,
|
get_attention_tp_rank,
|
||||||
get_attention_tp_size,
|
get_attention_tp_size,
|
||||||
|
get_dp_global_num_tokens,
|
||||||
get_global_dp_buffer,
|
get_global_dp_buffer,
|
||||||
get_local_dp_buffer,
|
get_local_dp_buffer,
|
||||||
is_allocation_symmetric,
|
is_allocation_symmetric,
|
||||||
@@ -55,6 +56,7 @@ from sglang.srt.layers.dp_attention import (
|
|||||||
from sglang.srt.layers.flashinfer_comm_fusion import is_flashinfer_allreduce_unavailable
|
from sglang.srt.layers.flashinfer_comm_fusion import is_flashinfer_allreduce_unavailable
|
||||||
from sglang.srt.layers.moe import (
|
from sglang.srt.layers.moe import (
|
||||||
get_moe_a2a_backend,
|
get_moe_a2a_backend,
|
||||||
|
should_use_dp_reduce_scatterv,
|
||||||
should_use_flashinfer_cutlass_moe_fp4_allgather,
|
should_use_flashinfer_cutlass_moe_fp4_allgather,
|
||||||
)
|
)
|
||||||
from sglang.srt.model_executor.forward_batch_info import ForwardBatch
|
from sglang.srt.model_executor.forward_batch_info import ForwardBatch
|
||||||
@@ -1096,8 +1098,13 @@ class CommunicateSummableTensorPairFn:
|
|||||||
get_local_dp_buffer(),
|
get_local_dp_buffer(),
|
||||||
hidden_states,
|
hidden_states,
|
||||||
)
|
)
|
||||||
if allow_reduce_scatter and forward_batch.dp_padding_mode.is_max_len():
|
if should_use_dp_reduce_scatterv():
|
||||||
# When using padding, all_reduce is skipped after MLP and MOE and reduce scatter is used here instead.
|
get_tp_group().reduce_scatterv(
|
||||||
|
global_hidden_states,
|
||||||
|
output=hidden_states,
|
||||||
|
sizes=get_dp_global_num_tokens(),
|
||||||
|
)
|
||||||
|
elif allow_reduce_scatter and forward_batch.dp_padding_mode.is_max_len():
|
||||||
dp_reduce_scatter_tensor(hidden_states, global_hidden_states)
|
dp_reduce_scatter_tensor(hidden_states, global_hidden_states)
|
||||||
else:
|
else:
|
||||||
dp_scatter(hidden_states, global_hidden_states, forward_batch)
|
dp_scatter(hidden_states, global_hidden_states, forward_batch)
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ from sglang.srt.layers.moe.utils import (
|
|||||||
get_tbo_token_distribution_threshold,
|
get_tbo_token_distribution_threshold,
|
||||||
initialize_moe_config,
|
initialize_moe_config,
|
||||||
is_tbo_enabled,
|
is_tbo_enabled,
|
||||||
|
should_use_dp_reduce_scatterv,
|
||||||
should_use_flashinfer_cutlass_moe_fp4_allgather,
|
should_use_flashinfer_cutlass_moe_fp4_allgather,
|
||||||
)
|
)
|
||||||
|
|
||||||
@@ -23,6 +24,7 @@ __all__ = [
|
|||||||
"get_moe_a2a_backend",
|
"get_moe_a2a_backend",
|
||||||
"get_moe_runner_backend",
|
"get_moe_runner_backend",
|
||||||
"get_deepep_mode",
|
"get_deepep_mode",
|
||||||
|
"should_use_dp_reduce_scatterv",
|
||||||
"should_use_flashinfer_cutlass_moe_fp4_allgather",
|
"should_use_flashinfer_cutlass_moe_fp4_allgather",
|
||||||
"is_tbo_enabled",
|
"is_tbo_enabled",
|
||||||
"get_tbo_token_distribution_threshold",
|
"get_tbo_token_distribution_threshold",
|
||||||
|
|||||||
@@ -298,6 +298,21 @@ def should_use_flashinfer_cutlass_moe_fp4_allgather():
|
|||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def should_use_dp_reduce_scatterv():
|
||||||
|
"""
|
||||||
|
Use reduce_scatterv in the standard dispatcher's combine() for DP attention
|
||||||
|
with EP, replacing the default all-reduce + dp_scatter path.
|
||||||
|
Only changes the combine (post-kernel) communication; dispatch is unchanged.
|
||||||
|
"""
|
||||||
|
return (
|
||||||
|
not should_use_flashinfer_cutlass_moe_fp4_allgather()
|
||||||
|
and get_moe_a2a_backend().is_none()
|
||||||
|
and is_dp_attention_enabled()
|
||||||
|
and get_attention_dp_size() > 1
|
||||||
|
and get_moe_expert_parallel_world_size() == get_attention_dp_size()
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
@contextmanager
|
@contextmanager
|
||||||
def speculative_moe_backend_context():
|
def speculative_moe_backend_context():
|
||||||
"""
|
"""
|
||||||
|
|||||||
@@ -57,6 +57,7 @@ from sglang.srt.layers.linear import (
|
|||||||
from sglang.srt.layers.logits_processor import LogitsProcessor
|
from sglang.srt.layers.logits_processor import LogitsProcessor
|
||||||
from sglang.srt.layers.moe import (
|
from sglang.srt.layers.moe import (
|
||||||
get_moe_a2a_backend,
|
get_moe_a2a_backend,
|
||||||
|
should_use_dp_reduce_scatterv,
|
||||||
should_use_flashinfer_cutlass_moe_fp4_allgather,
|
should_use_flashinfer_cutlass_moe_fp4_allgather,
|
||||||
)
|
)
|
||||||
from sglang.srt.layers.moe.ep_moe.layer import get_moe_impl_class
|
from sglang.srt.layers.moe.ep_moe.layer import get_moe_impl_class
|
||||||
@@ -352,6 +353,7 @@ class Qwen2MoeSparseMoeBlock(nn.Module):
|
|||||||
and not should_allreduce_fusion
|
and not should_allreduce_fusion
|
||||||
and not use_reduce_scatter
|
and not use_reduce_scatter
|
||||||
and not should_use_flashinfer_cutlass_moe_fp4_allgather()
|
and not should_use_flashinfer_cutlass_moe_fp4_allgather()
|
||||||
|
and not should_use_dp_reduce_scatterv()
|
||||||
):
|
):
|
||||||
final_hidden_states = tensor_model_parallel_all_reduce(final_hidden_states)
|
final_hidden_states = tensor_model_parallel_all_reduce(final_hidden_states)
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user