Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 28 additions & 1 deletion fastdeploy/spec_decode/mtp.py
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@
eagle_get_self_hidden_states,
mtp_save_first_token,
mtp_step_paddle,
set_data_ipc,
share_external_data,
)
from fastdeploy.model_executor.xpu_pre_and_post_process import (
Expand All @@ -65,6 +66,7 @@
speculate_get_logits,
speculate_save_output_topk,
update_attn_mask_offsets,
set_data_ipc,
)
from fastdeploy.model_executor.pre_and_post_process import pre_process, rebuild_padding

Expand Down Expand Up @@ -209,6 +211,9 @@ def initialize_kv_cache(self, main_model_num_blocks, profile: bool = False):
self.num_main_model_layers,
self.num_main_model_layers + self.model_config.num_hidden_layers,
):
logger.info(
f"..attaching kv cache for mtp layer {i}: key:{key_cache_shape}, value:{value_cache_shape}"
)
key_cache = paddle.empty(shape=[], dtype=cache_type)
key_cache_name = f"key_caches_{i}_rank{local_rank}.device{self.device_id}"
val_cache_name = f"value_caches_{i}_rank{local_rank}.device{self.device_id}"
Expand All @@ -232,28 +237,50 @@ def initialize_kv_cache(self, main_model_num_blocks, profile: bool = False):

self.model_inputs["caches"] = cache_kvs_list
else:
for i in range(self.model_config.num_hidden_layers):
for i in range(
self.num_main_model_layers,
self.num_main_model_layers + self.model_config.num_hidden_layers,
):
logger.info(f"..creating kv cache for mtp layer {i}: key:{key_cache_shape}, value:{value_cache_shape}")
self.cache_kvs[f"key_caches_{i}"] = paddle.full(
shape=key_cache_shape,
fill_value=0,
dtype=cache_type,
)
set_data_ipc(
self.cache_kvs[f"key_caches_{i}"], f"key_caches_{i}_rank{local_rank}.device{self.device_id}"
)

self.cache_kvs[f"value_caches_{i}"] = paddle.full(
shape=value_cache_shape,
fill_value=0,
dtype=cache_type,
)
set_data_ipc(
self.cache_kvs[f"value_caches_{i}"], f"value_caches_{i}_rank{local_rank}.device{self.device_id}"
)

if kv_cache_quant_type == "block_wise_fp8":
self.cache_kvs[f"key_cache_scales_{i}"] = paddle.full(
shape=kv_cache_scale_shape,
fill_value=0,
dtype=paddle.get_default_dtype(),
)
set_data_ipc(
self.cache_kvs[f"key_cache_scales_{i}"],
f"key_cache_scales_{i}_rank{local_rank}.device{self.device_id}",
)

self.cache_kvs[f"value_cache_scales_{i}"] = paddle.full(
shape=kv_cache_scale_shape,
fill_value=0,
dtype=paddle.get_default_dtype(),
)
set_data_ipc(
self.cache_kvs[f"value_cache_scales_{i}"],
f"value_cache_scales_{i}_rank{local_rank}.device{self.device_id}",
)

self.model_inputs["caches"] = list(self.cache_kvs.values())
for value in self.cache_kvs.values():
del value
Expand Down
2 changes: 2 additions & 0 deletions fastdeploy/worker/gpu_model_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -1664,10 +1664,12 @@ def initialize_kv_cache(self, profile: bool = False) -> None:
key_cache_scales = paddle.full(
shape=kv_cache_scale_shape, fill_value=0, dtype=paddle.get_default_dtype()
)
set_data_ipc(key_cache_scales, key_cache_scales_name)
if value_cache_shape:
val_cache_scales = paddle.full(
shape=kv_cache_scale_shape, fill_value=0, dtype=paddle.get_default_dtype()
)
set_data_ipc(val_cache_scales, value_cache_scales_name)
cache_kvs_list.extend([key_cache_scales, val_cache_scales])
else:
cache_kvs_list.extend([key_cache_scales])
Expand Down
6 changes: 2 additions & 4 deletions fastdeploy/worker/xpu_model_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -1499,10 +1499,9 @@ def profile_run(self) -> None:
"""Execute a forward pass with dummy inputs to profile the memory usage of the model"""

self.num_gpu_blocks = self.cache_config.total_block_num
self.initialize_kv_cache(profile=True)

if self.speculative_method in ["mtp"]:
self.proposer.initialize_kv_cache(main_model_num_blocks=self.num_gpu_blocks, profile=True)
self.initialize_kv_cache(profile=True)

self._dummy_run(
num_tokens=int(self.scheduler_config.max_num_batched_tokens),
Expand All @@ -1518,10 +1517,9 @@ def update_share_input_block_num(self, num_gpu_blocks: int) -> None:
self.num_gpu_blocks = num_gpu_blocks

# Reset block table and kv cache with global block num
self.initialize_kv_cache()

if self.speculative_method in ["mtp"]:
self.proposer.initialize_kv_cache(main_model_num_blocks=self.num_gpu_blocks)
self.initialize_kv_cache()

# Reset free list
free_list = list(
Expand Down
Loading