【NPU】Support EAGLE when PP enabled in prefill nodes (#32207)
This commit is contained in:
@@ -49,6 +49,25 @@ def check_server_args(server_args: Any):
|
|||||||
)
|
)
|
||||||
|
|
||||||
if cfg.pp_size > 1:
|
if cfg.pp_size > 1:
|
||||||
|
if get_platform().is_npu:
|
||||||
|
# NPU: allow PP + EAGLE speculative decoding
|
||||||
|
assert cfg.disable_overlap_schedule, (
|
||||||
|
"Pipeline parallelism is not compatible with overlap schedule"
|
||||||
|
)
|
||||||
|
if cfg.speculative_algorithm is not None:
|
||||||
|
assert (
|
||||||
|
cfg.speculative_algorithm.upper() == "EAGLE"
|
||||||
|
and not cfg.enable_multi_layer_eagle
|
||||||
|
), (
|
||||||
|
"Pipeline parallelism currently only supports EAGLE "
|
||||||
|
"(non-multi-layer) speculative decoding"
|
||||||
|
)
|
||||||
|
assert cfg.disaggregation_mode == "prefill", (
|
||||||
|
"NPU PP + speculative decoding (MTP) is only supported "
|
||||||
|
"on prefill nodes (disaggregation-mode=prefill)"
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
# Non-NPU: PP + speculative decoding is not supported
|
||||||
assert cfg.disable_overlap_schedule and cfg.speculative_algorithm is None, (
|
assert cfg.disable_overlap_schedule and cfg.speculative_algorithm is None, (
|
||||||
"Pipeline parallelism is not compatible with overlap schedule, speculative decoding"
|
"Pipeline parallelism is not compatible with overlap schedule, speculative decoding"
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -15,6 +15,7 @@ from sglang.srt.disaggregation.mooncake.conn import (
|
|||||||
MooncakeKVReceiver,
|
MooncakeKVReceiver,
|
||||||
MooncakeKVSender,
|
MooncakeKVSender,
|
||||||
)
|
)
|
||||||
|
from sglang.srt.distributed import get_pp_group
|
||||||
from sglang.srt.utils.network import get_local_ip_auto
|
from sglang.srt.utils.network import get_local_ip_auto
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
@@ -119,28 +120,37 @@ class AscendKVManager(MooncakeKVManager):
|
|||||||
# dst_kv_ptrs: k_data, v_data, index_k_data(optional)
|
# dst_kv_ptrs: k_data, v_data, index_k_data(optional)
|
||||||
# state_type is accepted for parity with the common disaggregation path;
|
# state_type is accepted for parity with the common disaggregation path;
|
||||||
# the NPU kv_buf_groups slicing below is state-type agnostic.
|
# the NPU kv_buf_groups slicing below is state-type agnostic.
|
||||||
start_layer = self.kv_args.prefill_start_layer
|
|
||||||
kv_buf_groups = getattr(self.kv_args, "kv_buf_groups", 1)
|
kv_buf_groups = getattr(self.kv_args, "kv_buf_groups", 1)
|
||||||
total_kv_layers = getattr(self.kv_args, "total_kv_layers", 0)
|
hidden_kv_layers = getattr(self.kv_args, "hidden_kv_layers", 0)
|
||||||
|
draft_kv_layers = getattr(self.kv_args, "draft_kv_layers", 0)
|
||||||
src_layers = len(src_kv_ptrs) // kv_buf_groups
|
src_layers = len(src_kv_ptrs) // kv_buf_groups
|
||||||
# When only speculative-algorithm is enabled for decode
|
dst_layers = len(dst_kv_ptrs) // kv_buf_groups
|
||||||
# the KV has one more layer than prefill.
|
if src_layers == dst_layers:
|
||||||
# The draft layer needs to be skipped.
|
|
||||||
dst_total_layers = (
|
|
||||||
min(len(dst_kv_ptrs) // kv_buf_groups, total_kv_layers)
|
|
||||||
if total_kv_layers
|
|
||||||
else len(dst_kv_ptrs) // kv_buf_groups
|
|
||||||
)
|
|
||||||
end_layer = start_layer + src_layers
|
|
||||||
if src_layers == dst_total_layers:
|
|
||||||
sliced_dst_kv_ptrs = dst_kv_ptrs
|
sliced_dst_kv_ptrs = dst_kv_ptrs
|
||||||
else:
|
else:
|
||||||
sliced_dst_kv_ptrs = []
|
sliced_dst_kv_ptrs = []
|
||||||
|
start_layer = self.kv_args.prefill_start_layer
|
||||||
|
transfer_draft_kv = get_pp_group().is_last_rank and draft_kv_layers
|
||||||
|
if transfer_draft_kv:
|
||||||
|
end_layer = start_layer + src_layers - draft_kv_layers
|
||||||
|
else:
|
||||||
|
end_layer = start_layer + src_layers
|
||||||
|
|
||||||
|
# target kv
|
||||||
for i in range(kv_buf_groups):
|
for i in range(kv_buf_groups):
|
||||||
layer_offset = i * dst_total_layers
|
layer_offset = i * hidden_kv_layers
|
||||||
sliced_dst_kv_ptrs.extend(
|
sliced_dst_kv_ptrs.extend(
|
||||||
dst_kv_ptrs[layer_offset + start_layer : layer_offset + end_layer]
|
dst_kv_ptrs[layer_offset + start_layer : layer_offset + end_layer]
|
||||||
)
|
)
|
||||||
|
# draft kv
|
||||||
|
if transfer_draft_kv:
|
||||||
|
for i in range(kv_buf_groups):
|
||||||
|
layer_offset = (
|
||||||
|
i * draft_kv_layers + kv_buf_groups * hidden_kv_layers
|
||||||
|
)
|
||||||
|
sliced_dst_kv_ptrs.extend(
|
||||||
|
dst_kv_ptrs[layer_offset : layer_offset + draft_kv_layers]
|
||||||
|
)
|
||||||
layers_current_pp_stage = len(src_kv_ptrs)
|
layers_current_pp_stage = len(src_kv_ptrs)
|
||||||
return src_kv_ptrs, sliced_dst_kv_ptrs, layers_current_pp_stage
|
return src_kv_ptrs, sliced_dst_kv_ptrs, layers_current_pp_stage
|
||||||
|
|
||||||
|
|||||||
@@ -92,7 +92,9 @@ class KVArgs:
|
|||||||
# Only used of npu, for kv buf groups
|
# Only used of npu, for kv buf groups
|
||||||
kv_buf_groups: int
|
kv_buf_groups: int
|
||||||
# Only used of npu, for decode total kv layers
|
# Only used of npu, for decode total kv layers
|
||||||
total_kv_layers: int
|
hidden_kv_layers: int
|
||||||
|
# Only used of npu, for decode total kv layers
|
||||||
|
draft_kv_layers: int
|
||||||
|
|
||||||
|
|
||||||
class KVPoll:
|
class KVPoll:
|
||||||
|
|||||||
@@ -56,6 +56,7 @@ from sglang.srt.disaggregation.utils import (
|
|||||||
prepare_abort,
|
prepare_abort,
|
||||||
setup_state_kv_args,
|
setup_state_kv_args,
|
||||||
)
|
)
|
||||||
|
from sglang.srt.distributed import get_pp_group
|
||||||
from sglang.srt.environ import envs
|
from sglang.srt.environ import envs
|
||||||
from sglang.srt.managers.schedule_batch import (
|
from sglang.srt.managers.schedule_batch import (
|
||||||
FINISH_ABORT,
|
FINISH_ABORT,
|
||||||
@@ -246,7 +247,11 @@ class PrefillBootstrapQueue:
|
|||||||
else getattr(self.token_to_kv_pool, "end_layer", None)
|
else getattr(self.token_to_kv_pool, "end_layer", None)
|
||||||
)
|
)
|
||||||
|
|
||||||
draft_kv_pool = self.draft_token_to_kv_pool if transfer_draft_cache else None
|
draft_kv_pool = (
|
||||||
|
self.draft_token_to_kv_pool
|
||||||
|
if transfer_draft_cache and (not _is_npu or get_pp_group().is_last_rank)
|
||||||
|
else None
|
||||||
|
)
|
||||||
num_draft_entries = 0
|
num_draft_entries = 0
|
||||||
if draft_kv_pool is not None:
|
if draft_kv_pool is not None:
|
||||||
# We should also transfer draft model kv cache. The indices are
|
# We should also transfer draft model kv cache. The indices are
|
||||||
|
|||||||
@@ -1466,7 +1466,10 @@ def setup_state_kv_args(
|
|||||||
kv_args.kv_buf_groups = (
|
kv_args.kv_buf_groups = (
|
||||||
len(kv_args.kv_data_ptrs) // token_to_kv_pool.layer_num
|
len(kv_args.kv_data_ptrs) // token_to_kv_pool.layer_num
|
||||||
)
|
)
|
||||||
kv_args.total_kv_layers = total_kv_layers
|
kv_args.hidden_kv_layers = total_kv_layers
|
||||||
|
kv_args.draft_kv_layers = (
|
||||||
|
draft_token_to_kv_pool.layer_num if draft_token_to_kv_pool else 0
|
||||||
|
)
|
||||||
else:
|
else:
|
||||||
append_state_component(
|
append_state_component(
|
||||||
kv_args, StateType.DSA, data_ptrs, data_lens, item_lens
|
kv_args, StateType.DSA, data_ptrs, data_lens, item_lens
|
||||||
|
|||||||
@@ -447,7 +447,7 @@ def forward_dsa_prepare_npu(
|
|||||||
|
|
||||||
q_nope_out = q_nope_out.transpose(0, 1)
|
q_nope_out = q_nope_out.transpose(0, 1)
|
||||||
|
|
||||||
if m.layer_id == 0:
|
if m.layer_id == get_token_to_kv_pool().start_layer:
|
||||||
m.rotary_emb.sin_cos_cache = m.rotary_emb.cos_sin_cache.index_select(
|
m.rotary_emb.sin_cos_cache = m.rotary_emb.cos_sin_cache.index_select(
|
||||||
0, positions
|
0, positions
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -159,7 +159,7 @@ class DSANPUIndexerMixin:
|
|||||||
|
|
||||||
k_pe = k_pe.unsqueeze(1)
|
k_pe = k_pe.unsqueeze(1)
|
||||||
|
|
||||||
if layer_id == 0:
|
if layer_id == get_token_to_kv_pool().start_layer:
|
||||||
self.rotary_emb.sin_cos_cache = (
|
self.rotary_emb.sin_cos_cache = (
|
||||||
self.rotary_emb.cos_sin_cache.index_select(0, positions)
|
self.rotary_emb.cos_sin_cache.index_select(0, positions)
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -608,7 +608,6 @@ class FutureMap:
|
|||||||
self.output_tokens_buf[indices] = payload.bonus_tokens.to(
|
self.output_tokens_buf[indices] = payload.bonus_tokens.to(
|
||||||
self.output_tokens_buf.dtype
|
self.output_tokens_buf.dtype
|
||||||
)
|
)
|
||||||
|
|
||||||
if self.need_topk:
|
if self.need_topk:
|
||||||
self.topk_p_buf[indices] = payload.topk_p.to(self.topk_p_buf.dtype)
|
self.topk_p_buf[indices] = payload.topk_p.to(self.topk_p_buf.dtype)
|
||||||
self.topk_index_buf[indices] = payload.topk_index.to(
|
self.topk_index_buf[indices] = payload.topk_index.to(
|
||||||
|
|||||||
@@ -5,6 +5,8 @@ from typing import TYPE_CHECKING, Any, NamedTuple
|
|||||||
import msgspec
|
import msgspec
|
||||||
from torch import nn
|
from torch import nn
|
||||||
|
|
||||||
|
from sglang.srt.utils import is_npu
|
||||||
|
|
||||||
if TYPE_CHECKING:
|
if TYPE_CHECKING:
|
||||||
from sglang.srt.configs.model_config import ModelConfig
|
from sglang.srt.configs.model_config import ModelConfig
|
||||||
from sglang.srt.speculative.spec_info import SpeculativeAlgorithm
|
from sglang.srt.speculative.spec_info import SpeculativeAlgorithm
|
||||||
@@ -156,6 +158,7 @@ def resolve_layer_indices(
|
|||||||
if loop_num > 1:
|
if loop_num > 1:
|
||||||
num_effective_layers = num_effective_layers * loop_num
|
num_effective_layers = num_effective_layers * loop_num
|
||||||
|
|
||||||
|
if not is_npu():
|
||||||
_assert_pp_mtp_compat(
|
_assert_pp_mtp_compat(
|
||||||
model_has_mtp_layers=model_has_mtp_layers,
|
model_has_mtp_layers=model_has_mtp_layers,
|
||||||
spec_algorithm=spec_algorithm,
|
spec_algorithm=spec_algorithm,
|
||||||
|
|||||||
@@ -39,7 +39,7 @@ from sglang.srt.layers.vocab_parallel_embedding import (
|
|||||||
VocabParallelEmbedding,
|
VocabParallelEmbedding,
|
||||||
get_embedding_tp_kwargs,
|
get_embedding_tp_kwargs,
|
||||||
)
|
)
|
||||||
from sglang.srt.model_executor.forward_batch_info import ForwardBatch
|
from sglang.srt.model_executor.forward_batch_info import ForwardBatch, PPProxyTensors
|
||||||
from sglang.srt.models.deepseek_common.utils import enable_nextn_moe_bf16_cast_to_fp8
|
from sglang.srt.models.deepseek_common.utils import enable_nextn_moe_bf16_cast_to_fp8
|
||||||
from sglang.srt.models.deepseek_v2 import DeepseekV2DecoderLayer, DeepseekV3ForCausalLM
|
from sglang.srt.models.deepseek_v2 import DeepseekV2DecoderLayer, DeepseekV3ForCausalLM
|
||||||
from sglang.srt.models.utils import WeightsMapper
|
from sglang.srt.models.utils import WeightsMapper
|
||||||
@@ -303,6 +303,7 @@ class DeepseekV3ForCausalLMNextN(DeepseekV3ForCausalLM):
|
|||||||
input_ids: torch.Tensor,
|
input_ids: torch.Tensor,
|
||||||
positions: torch.Tensor,
|
positions: torch.Tensor,
|
||||||
forward_batch: ForwardBatch,
|
forward_batch: ForwardBatch,
|
||||||
|
pp_proxy_tensors: Optional[PPProxyTensors] = None,
|
||||||
) -> torch.Tensor:
|
) -> torch.Tensor:
|
||||||
hidden_states = self.model(input_ids, positions, forward_batch)
|
hidden_states = self.model(input_ids, positions, forward_batch)
|
||||||
return self.logits_processor(
|
return self.logits_processor(
|
||||||
|
|||||||
@@ -2588,7 +2588,7 @@ class DeepseekV2Model(nn.Module):
|
|||||||
self.first_k_dense_replace = config.first_k_dense_replace
|
self.first_k_dense_replace = config.first_k_dense_replace
|
||||||
self.pp_group = get_pp_group()
|
self.pp_group = get_pp_group()
|
||||||
|
|
||||||
if self.pp_group.is_first_rank:
|
if self.pp_group.is_first_rank or (_is_npu and self.pp_group.is_last_rank):
|
||||||
self.embed_tokens = VocabParallelEmbedding(
|
self.embed_tokens = VocabParallelEmbedding(
|
||||||
config.vocab_size,
|
config.vocab_size,
|
||||||
config.hidden_size,
|
config.hidden_size,
|
||||||
|
|||||||
@@ -41,7 +41,11 @@ from sglang.srt.model_executor.cuda_graph_config import (
|
|||||||
Phase,
|
Phase,
|
||||||
check_cuda_graph_backend,
|
check_cuda_graph_backend,
|
||||||
)
|
)
|
||||||
from sglang.srt.model_executor.forward_batch_info import CaptureHiddenMode, ForwardBatch
|
from sglang.srt.model_executor.forward_batch_info import (
|
||||||
|
CaptureHiddenMode,
|
||||||
|
ForwardBatch,
|
||||||
|
PPProxyTensors,
|
||||||
|
)
|
||||||
from sglang.srt.model_executor.forward_context import ForwardContext, forward_context
|
from sglang.srt.model_executor.forward_context import ForwardContext, forward_context
|
||||||
from sglang.srt.model_executor.runner import (
|
from sglang.srt.model_executor.runner import (
|
||||||
DecodeCudaGraphRunner,
|
DecodeCudaGraphRunner,
|
||||||
@@ -1190,7 +1194,7 @@ class EAGLEWorkerV2(BaseSpecWorker):
|
|||||||
batch: ScheduleBatch,
|
batch: ScheduleBatch,
|
||||||
on_publish=None,
|
on_publish=None,
|
||||||
grammar_barrier=None,
|
grammar_barrier=None,
|
||||||
pp_proxy_tensors=None,
|
pp_proxy_tensors: Optional[PPProxyTensors] = None,
|
||||||
):
|
):
|
||||||
if batch.forward_mode.is_extend() or batch.is_extend_in_batch:
|
if batch.forward_mode.is_extend() or batch.is_extend_in_batch:
|
||||||
# Target prefill
|
# Target prefill
|
||||||
|
|||||||
Reference in New Issue
Block a user