Use reduce scatter for DP (#8539)

This commit is contained in:
Trevor Morris
2025-08-06 16:21:26 -07:00
committed by GitHub
parent 92cc32d9fc
commit c0e84297c2
6 changed files with 73 additions and 18 deletions
+33 -9
View File
@@ -208,13 +208,21 @@ class DeepseekV2MLP(nn.Module):
)
self.act_fn = SiluAndMul()
def forward(self, x, forward_batch=None, can_fuse_mlp_allreduce=False):
def forward(
self,
x,
forward_batch=None,
can_fuse_mlp_allreduce: bool = False,
use_reduce_scatter: bool = False,
):
if (self.tp_size == 1) and x.shape[0] == 0:
return x
gate_up, _ = self.gate_up_proj(x)
x = self.act_fn(gate_up)
x, _ = self.down_proj(x, can_fuse_mlp_allreduce=can_fuse_mlp_allreduce)
x, _ = self.down_proj(
x, skip_all_reduce=can_fuse_mlp_allreduce or use_reduce_scatter
)
return x
@@ -441,6 +449,7 @@ class DeepseekV2MoE(nn.Module):
hidden_states: torch.Tensor,
forward_batch: Optional[ForwardBatch] = None,
can_fuse_mlp_allreduce: bool = False,
use_reduce_scatter: bool = False,
) -> torch.Tensor:
if not self._enable_deepep_moe:
DUAL_STREAM_TOKEN_THRESHOLD = 1024
@@ -450,15 +459,20 @@ class DeepseekV2MoE(nn.Module):
and hidden_states.shape[0] <= DUAL_STREAM_TOKEN_THRESHOLD
):
return self.forward_normal_dual_stream(
hidden_states, can_fuse_mlp_allreduce
hidden_states, can_fuse_mlp_allreduce, use_reduce_scatter
)
else:
return self.forward_normal(hidden_states, can_fuse_mlp_allreduce)
return self.forward_normal(
hidden_states, can_fuse_mlp_allreduce, use_reduce_scatter
)
else:
return self.forward_deepep(hidden_states, forward_batch)
def forward_normal_dual_stream(
self, hidden_states: torch.Tensor, can_fuse_mlp_allreduce: bool = False
self,
hidden_states: torch.Tensor,
can_fuse_mlp_allreduce: bool = False,
use_reduce_scatter: bool = False,
) -> torch.Tensor:
current_stream = torch.cuda.current_stream()
@@ -486,12 +500,15 @@ class DeepseekV2MoE(nn.Module):
torch.add(final_hidden_states, shared_output, out=final_hidden_states_out)
final_hidden_states = final_hidden_states_out
sm.tag(final_hidden_states)
if self.tp_size > 1 and not can_fuse_mlp_allreduce:
if self.tp_size > 1 and not can_fuse_mlp_allreduce and not use_reduce_scatter:
final_hidden_states = tensor_model_parallel_all_reduce(final_hidden_states)
return final_hidden_states
def forward_normal(
self, hidden_states: torch.Tensor, can_fuse_mlp_allreduce: bool = False
self,
hidden_states: torch.Tensor,
can_fuse_mlp_allreduce: bool = False,
use_reduce_scatter: bool = False,
) -> torch.Tensor:
if hasattr(self, "shared_experts") and use_intel_amx_backend(
self.shared_experts.gate_up_proj
@@ -520,7 +537,7 @@ class DeepseekV2MoE(nn.Module):
torch.add(final_hidden_states, shared_output, out=final_hidden_states_out)
final_hidden_states = final_hidden_states_out
sm.tag(final_hidden_states)
if self.tp_size > 1 and not can_fuse_mlp_allreduce:
if self.tp_size > 1 and not can_fuse_mlp_allreduce and not use_reduce_scatter:
final_hidden_states = tensor_model_parallel_all_reduce(final_hidden_states)
return final_hidden_states
@@ -1822,6 +1839,7 @@ class DeepseekV2DecoderLayer(nn.Module):
layer_scatter_modes=self.layer_scatter_modes,
input_layernorm=self.input_layernorm,
post_attention_layernorm=self.post_attention_layernorm,
allow_reduce_scatter=True,
)
def _is_layer_sparse(self, layer_id: int, is_nextn: bool) -> bool:
@@ -1884,7 +1902,13 @@ class DeepseekV2DecoderLayer(nn.Module):
and not self.is_nextn
)
hidden_states = self.mlp(hidden_states, forward_batch, can_fuse_mlp_allreduce)
# For DP with padding, reduce scatter can be used instead of all-reduce.
use_reduce_scatter = self.layer_communicator.should_use_reduce_scatter(
forward_batch
)
hidden_states = self.mlp(
hidden_states, forward_batch, can_fuse_mlp_allreduce, use_reduce_scatter
)
if can_fuse_mlp_allreduce:
hidden_states._sglang_needs_allreduce_fusion = True
+1 -1
View File
@@ -160,7 +160,7 @@ class Glm4MoeMLP(nn.Module):
gate_up, _ = self.gate_up_proj(x)
x = self.act_fn(gate_up)
x, _ = self.down_proj(x, can_fuse_mlp_allreduce=can_fuse_mlp_allreduce)
x, _ = self.down_proj(x, skip_all_reduce=can_fuse_mlp_allreduce)
return x