From d20d9b0b7b81f9498fb296c751fe1789cc98e9cb Mon Sep 17 00:00:00 2001 From: ltd0924 Date: Mon, 1 Sep 2025 10:52:32 +0800 Subject: [PATCH 1/8] support model weight update in ep --- fastdeploy/model_executor/layers/moe/ep.py | 11 ++++++- .../layers/moe/fused_moe_backend_base.py | 4 +++ fastdeploy/rl/dynamic_weight_manager.py | 16 +++++---- fastdeploy/worker/worker_process.py | 33 +++++++++++-------- 4 files changed, 44 insertions(+), 20 deletions(-) diff --git a/fastdeploy/model_executor/layers/moe/ep.py b/fastdeploy/model_executor/layers/moe/ep.py index 9659aec7d47..3c2ac8e64d5 100644 --- a/fastdeploy/model_executor/layers/moe/ep.py +++ b/fastdeploy/model_executor/layers/moe/ep.py @@ -78,6 +78,7 @@ def __init__( splitwise_role: str, moe_phase: MoEPhase, async_finish: bool = False, + group=None, ): """ Initialize the DeepEP engine. @@ -90,7 +91,9 @@ def __init__( num_experts: The number of experts. """ # TODO(@wufeisheng): Support configurable EP size​ - self.group = paddle.distributed.new_group(range(ep_size)) + if group is None: + group = paddle.distributed.new_group(range(ep_size)) + self.group = group self.ep_size = ep_size self.rank_id = ep_rank self.hidden = hidden @@ -277,6 +280,7 @@ def __init__( ep_size: int = 1, ep_rank: int = 0, redundant_experts_num: int = 0, + ep_group=None, ): self.top_k = top_k self.num_experts = num_experts @@ -289,6 +293,7 @@ def __init__( ep_rank=ep_rank, splitwise_role=splitwise_role, moe_phase=moe_phase, + group=ep_group, ) def moe_select(self, layer: nn.Layer, gate_out: paddle.Tensor): @@ -367,6 +372,7 @@ def __init__( ep_size: int = 1, ep_rank: int = 0, redundant_experts_num: int = 0, + ep_group=None, moe_phase: MoEPhase = MoEPhase("prefill"), ): super().__init__( @@ -379,6 +385,7 @@ def __init__( ep_size=ep_size, ep_rank=ep_rank, redundant_experts_num=redundant_experts_num, + ep_group=ep_group, ) def dispatch( @@ -445,6 +452,7 @@ def __init__( ep_size: int = 1, ep_rank: int = 0, redundant_experts_num: int = 0, + ep_group=None, moe_phase: MoEPhase = MoEPhase("decode"), ): super().__init__( @@ -457,6 +465,7 @@ def __init__( ep_size=ep_size, ep_rank=ep_rank, redundant_experts_num=redundant_experts_num, + ep_group=ep_group, ) def dispatch( diff --git a/fastdeploy/model_executor/layers/moe/fused_moe_backend_base.py b/fastdeploy/model_executor/layers/moe/fused_moe_backend_base.py index bb29f3e3a3a..82df14a36b4 100644 --- a/fastdeploy/model_executor/layers/moe/fused_moe_backend_base.py +++ b/fastdeploy/model_executor/layers/moe/fused_moe_backend_base.py @@ -58,6 +58,7 @@ def init_ep(self, layer: nn.Layer) -> None: layer.ep_size, layer.ep_rank, layer.fd_config.model_config.redundant_experts_num, + layer.fd_config.parallel_config.ep_group, ) self.ep_decoder_runner = EPDecoderRunner( layer.top_k, @@ -68,6 +69,7 @@ def init_ep(self, layer: nn.Layer) -> None: layer.ep_size, layer.ep_rank, layer.fd_config.model_config.redundant_experts_num, + layer.fd_config.parallel_config.ep_group, ) else: if layer.fd_config.parallel_config.moe_phase.phase == "prefill": @@ -82,6 +84,7 @@ def init_ep(self, layer: nn.Layer) -> None: layer.ep_size, layer.ep_rank, layer.fd_config.model_config.redundant_experts_num, + layer.fd_config.parallel_config.ep_group, ) else: from .ep import EPDecoderRunner @@ -95,6 +98,7 @@ def init_ep(self, layer: nn.Layer) -> None: layer.ep_size, layer.ep_rank, layer.fd_config.model_config.redundant_experts_num, + layer.fd_config.parallel_config.ep_group, ) def process_loaded_weights(self, layer, weights) -> None: diff --git a/fastdeploy/rl/dynamic_weight_manager.py b/fastdeploy/rl/dynamic_weight_manager.py index ad39accdb2b..d00db929140 100644 --- a/fastdeploy/rl/dynamic_weight_manager.py +++ b/fastdeploy/rl/dynamic_weight_manager.py @@ -63,7 +63,9 @@ def update_parameters(self, pid: int = 0) -> None: paddle.device.cuda.empty_cache() if not self.first_load: - paddle.distributed.restart_process_group() + paddle.distributed.restart_process_group(self.parallel_config.tp_group) + if self.parallel_config.enable_expert_parallel: + paddle.distributed.restart_process_group(self.parallel_config.ep_group) strategy_handlers = { "ipc_snapshot": self._update_ipc_snapshot, @@ -110,9 +112,11 @@ def clear_parameters(self, pid: int = 0) -> None: param._clear_data() self._verify_parameters("clearance") - if self.nranks > 1: - paddle.distributed.barrier() - paddle.distributed.shutdown_process_group() + if self.parallel_config.tensor_parallel_size > 1: + paddle.distributed.barrier(self.parallel_config.tp_group) + paddle.distributed.shutdown_process_group(self.parallel_config.tp_group) + if self.parallel_config.enable_expert_parallel: + paddle.distributed.barrier(self.parallel_config.ep_group) self._update_shared_status(pid, -2) def _update_model_from_state(self, state_dict: Dict[str, paddle.Tensor], src_type: str): @@ -141,8 +145,8 @@ def _validate_parameter_match(self, name: str, src: paddle.Tensor, dst: paddle.T def _finalize_update(self, pid: int): """Finalize update process with verification.""" self._verify_parameters("update") - if self.nranks > 1: - paddle.distributed.barrier() + if self.parallel_config.tensor_parallel_size > 1: + paddle.distributed.barrier(self.parallel_config.tp_group) if not self.first_load: self._update_shared_status(pid, 0) self.first_load = False diff --git a/fastdeploy/worker/worker_process.py b/fastdeploy/worker/worker_process.py index d2ab0db9955..bdf37f6ffad 100644 --- a/fastdeploy/worker/worker_process.py +++ b/fastdeploy/worker/worker_process.py @@ -254,16 +254,15 @@ def event_loop_normal(self) -> None: """ # Currently, only support single node self.nnode = int((self.parallel_config.tensor_parallel_size + 7) // 8) - mp_num_per_node = self.parallel_config.tensor_parallel_size // self.nnode req_ids = [] num_running_requests = 0 - local_rank = self.local_rank % self.parallel_config.tensor_parallel_size + + self.model_weights_signal = paddle.zeros([1], dtype=paddle.int32) while True: - if self.local_rank == 0: + if self.local_rank % self.parallel_config.tensor_parallel_size == 0: if self.model_weights_status.value[0] != 0: - self.exist_task_signal.value[0] = 2 - else: - self.exist_task_signal.value[0] = 0 + self.model_weights_signal[0] = int(self.model_weights_status.value[0]) + paddle.distributed.broadcast(self.model_weights_signal, src=0) if self.parallel_config.tensor_parallel_size > 1: # Synchronize before updating weights @@ -271,10 +270,11 @@ def event_loop_normal(self) -> None: self.insert_step = False req_dicts = None + local_rank = self.local_rank % self.parallel_config.tensor_parallel_size self.worker_healthy_live_signal.value[local_rank % self.max_chips_per_node] = int(time.time()) # The first worker detects whether there are tasks in the task queue - if self.local_rank % mp_num_per_node == 0: + if self.local_rank % self.parallel_config.tensor_parallel_size == 0: if self.task_queue.num_tasks() > 0: # VL only support 1 batch to prefill if envs.ENABLE_V1_KVCACHE_SCHEDULER or not ( @@ -285,21 +285,28 @@ def event_loop_normal(self) -> None: else: self.exist_task_signal.value[0] = 1 - if self.parallel_config.tensor_parallel_size > 1: - # Synchronize the signal for other workers - paddle.distributed.barrier(self.parallel_config.tp_group) - if self.fd_config.load_config.dynamic_load_weight: - if self.exist_task_signal.value[0] == 2: + if self.parallel_config.enable_expert_parallel: + paddle.distributed.barrier(group=self.parallel_config.ep_group) + else: + paddle.distributed.barrier(self.parallel_config.tp_group) + if self.model_weights_signal[0] != 0: + logger.info(f"Rank: {self.local_rank} has updated parameters.") from fastdeploy.rl.dynamic_weight_manager import ( DynamicWeightManager, ) + self.model_weights_status.value[0] = self.model_weights_signal[0] + paddle.distributed.barrier(self.parallel_config.ep_group) + DynamicWeightManager.check_model_weights_status( self.model_weights_status, + # model_weights_signal self.worker.model_runner, - self.parallel_config.engine_pid, + self.parallel_config.engine_worker_queue_port, ) + self.model_weights_signal[0] = 0 + paddle.distributed.broadcast(self.model_weights_signal, src=0) if self.exist_task_signal.value[0] == 1 or self.task_queue.read_finish_flag.get() == 1: logger.info(f"Rank: {self.local_rank} Detected new requests.") From ba317d0a9dbb1fd09b38f143fd23ab4aa78192e1 Mon Sep 17 00:00:00 2001 From: ltd0924 Date: Mon, 1 Sep 2025 15:08:04 +0800 Subject: [PATCH 2/8] support model weight update in ep --- fastdeploy/config.py | 75 ++------------------------------------------ 1 file changed, 2 insertions(+), 73 deletions(-) diff --git a/fastdeploy/config.py b/fastdeploy/config.py index e4182e6c9a7..bb3e4900046 100644 --- a/fastdeploy/config.py +++ b/fastdeploy/config.py @@ -350,8 +350,8 @@ def set_tp_group(self): ) ) # same ep group id - # (TODO:gaoziyuan move this gid config to ep.py) dist.collective._set_custom_gid(self.data_parallel_size + tp_gid_offset) + self.ep_group = dist.new_group(range(self.expert_parallel_size)) logger.info( f"data_parallel_size: {self.data_parallel_size}, tensor_parallel_size: {self.tensor_parallel_size}, expert_parallel_size: {self.expert_parallel_size}, data_parallel_rank: {self.data_parallel_rank}, tensor_parallel_rank: {self.tensor_parallel_rank}, expert_parallel_rank: {self.expert_parallel_rank}, tp_group: {self.tp_group}." ) @@ -684,67 +684,6 @@ def update_use_cudagraph(self, argument: bool): argument = self.use_cudagraph -class MobaAttentionConfig: - def __init__( - self, - args, - ): - self.moba_encoder_top_k_left: int = None - self.moba_encoder_top_k_right: int = None - "The sparse topk of encoder attention is located at [moba_encoder_top_k_left, moba_encoder top_k_right]" - self.moba_decoder_top_k_left: int = None - self.moba_decoder_top_k_right: int = None - "The sparse topk of decoder attention is located at [moba_decoder_top_k_left, moba_decoder top_k_right]" - self.moba_use_encoder_seq_limit: int = None - "When the number of encdoer token is less than moba_use_encoder_seq_limit, it is not sparse" - self.moba_use_decoder_seq_limit: int = None - "When the number of decdoer token is less than moba_use_decoder_seq_limit, it is not sparse" - self.moba_block_size: int = 128 - self.mlp_weight_name: str = "moba_mlp_weight.safetensors" - self.moba_max_seq_length: int = 128 * 1024 - if args is not None: - for key, value in args.items(): - if hasattr(self, key): - setattr(self, key, value) - if self.moba_use_encoder_seq_limit is None and self.moba_encoder_top_k_left is not None: - self.moba_use_encoder_seq_limit = self.moba_encoder_top_k_left * self.moba_block_size - if self.moba_use_decoder_seq_limit is None and self.moba_decoder_top_k_left is not None: - self.moba_use_decoder_seq_limit = self.moba_decoder_top_k_left * self.moba_block_size - self.check_legality_parameters() - - def check_legality_parameters( - self, - ) -> None: - if self.moba_encoder_top_k_left is not None: - assert self.moba_encoder_top_k_left > 0, "moba_encoder_top_k_left must large than 0" - - if self.moba_encoder_top_k_right is not None: - assert self.moba_encoder_top_k_right > 0, "moba_encoder_top_k_right must large than 0" - assert ( - self.moba_encoder_top_k_right >= self.moba_encoder_top_k_left - ), "moba_encoder_top_k_right must large than moba_encoder_top_k_left" - - if self.moba_decoder_top_k_left is not None: - assert self.moba_decoder_top_k_left > 0, "moba_decoder_top_k_left must large than 0" - - if self.moba_decoder_top_k_right is not None: - assert self.moba_decoder_top_k_right > 0, "moba_decoder_top_k_right must large than 0" - assert ( - self.moba_decoder_top_k_right >= self.moba_decoder_top_k_left - ), "moba_decoder_top_k_right must large than moba_decoder_top_k_left" - - if self.moba_use_encoder_seq_limit is not None and self.moba_encoder_top_k_left is not None: - assert self.moba_use_encoder_seq_limit >= self.moba_encoder_top_k_left * self.moba_block_size - if self.moba_use_decoder_seq_limit is not None and self.moba_decoder_top_k_left is not None: - assert self.moba_use_decoder_seq_limit >= self.moba_decoder_top_k_left * self.moba_block_size - - def to_json_string(self): - """ - Convert moba_attention_config to json string. - """ - return json.dumps({key: value for key, value in self.__dict__.items() if value is not None}) - - class EarlyStopConfig: def __init__( self, @@ -1099,7 +1038,6 @@ def __init__( decoding_config: DecodingConfig = None, quant_config: QuantConfigBase = None, graph_opt_config: GraphOptimizationConfig = None, - moba_attention_config: MobaAttentionConfig = None, speculative_config: SpeculativeConfig = None, tokenizer: str = None, max_model_len: int = 8192, @@ -1134,7 +1072,7 @@ def __init__( self.early_stop_config: Optional[EarlyStopConfig] = early_stop_config self.decoding_config: DecodingConfig = decoding_config # type: ignore self.cache_config: CacheConfig = cache_config # type: ignore - self.moba_attention_config: Optional[MobaAttentionConfig] = moba_attention_config + # Initialize cuda graph capture list if self.graph_opt_config.cudagraph_capture_sizes is None: self.graph_opt_config._set_cudagraph_sizes(max_num_seqs=self.parallel_config.max_num_seqs) @@ -1233,15 +1171,6 @@ def postprocess(self): self.paddle_commit_id = paddle.version.commit - if self.cache_config.enable_chunked_prefill: - self.force_chunked_prefill = int(envs.FD_FORCE_CHUNKED_PREFILL) - if ( - self.speculative_config is not None - and self.speculative_config.method in ["mtp"] - and not self.force_chunked_prefill - ): - self.cache_config.enable_chunked_prefill = False - if self.max_num_batched_tokens is None: if self.cache_config.enable_chunked_prefill: self.max_num_batched_tokens = 2048 From 8b1212df096f7c17ae6b93c9b80293513217addd Mon Sep 17 00:00:00 2001 From: ltd0924 Date: Mon, 1 Sep 2025 17:01:36 +0800 Subject: [PATCH 3/8] support model weight update in ep --- fastdeploy/worker/worker_process.py | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/fastdeploy/worker/worker_process.py b/fastdeploy/worker/worker_process.py index bdf37f6ffad..3ee3e2647d1 100644 --- a/fastdeploy/worker/worker_process.py +++ b/fastdeploy/worker/worker_process.py @@ -262,7 +262,10 @@ def event_loop_normal(self) -> None: if self.local_rank % self.parallel_config.tensor_parallel_size == 0: if self.model_weights_status.value[0] != 0: self.model_weights_signal[0] = int(self.model_weights_status.value[0]) - paddle.distributed.broadcast(self.model_weights_signal, src=0) + if self.fd_config.load_config.dynamic_load_weight and self.parallel_config.enable_expert_parallel: + paddle.distributed.broadcast(self.model_weights_signal, src=0, group=self.parallel_config.ep_group) + if self.fd_config.load_config.dynamic_load_weight: + paddle.distributed.broadcast(self.model_weights_signal, src=0, group=self.parallel_config.tp_group) if self.parallel_config.tensor_parallel_size > 1: # Synchronize before updating weights @@ -297,8 +300,6 @@ def event_loop_normal(self) -> None: ) self.model_weights_status.value[0] = self.model_weights_signal[0] - paddle.distributed.barrier(self.parallel_config.ep_group) - DynamicWeightManager.check_model_weights_status( self.model_weights_status, # model_weights_signal From 6b80649c1816e1b3e1918a0ff4710bef61cfb5ad Mon Sep 17 00:00:00 2001 From: ltd0924 Date: Mon, 1 Sep 2025 17:05:33 +0800 Subject: [PATCH 4/8] support model weight update in ep --- fastdeploy/config.py | 73 +++++++++++++++++++++++++++++++++++++++++++- 1 file changed, 72 insertions(+), 1 deletion(-) diff --git a/fastdeploy/config.py b/fastdeploy/config.py index bb3e4900046..09e5c684bd3 100644 --- a/fastdeploy/config.py +++ b/fastdeploy/config.py @@ -684,6 +684,67 @@ def update_use_cudagraph(self, argument: bool): argument = self.use_cudagraph +class MobaAttentionConfig: + def __init__( + self, + args, + ): + self.moba_encoder_top_k_left: int = None + self.moba_encoder_top_k_right: int = None + "The sparse topk of encoder attention is located at [moba_encoder_top_k_left, moba_encoder top_k_right]" + self.moba_decoder_top_k_left: int = None + self.moba_decoder_top_k_right: int = None + "The sparse topk of decoder attention is located at [moba_decoder_top_k_left, moba_decoder top_k_right]" + self.moba_use_encoder_seq_limit: int = None + "When the number of encdoer token is less than moba_use_encoder_seq_limit, it is not sparse" + self.moba_use_decoder_seq_limit: int = None + "When the number of decdoer token is less than moba_use_decoder_seq_limit, it is not sparse" + self.moba_block_size: int = 128 + self.mlp_weight_name: str = "moba_mlp_weight.safetensors" + self.moba_max_seq_length: int = 128 * 1024 + if args is not None: + for key, value in args.items(): + if hasattr(self, key): + setattr(self, key, value) + if self.moba_use_encoder_seq_limit is None and self.moba_encoder_top_k_left is not None: + self.moba_use_encoder_seq_limit = self.moba_encoder_top_k_left * self.moba_block_size + if self.moba_use_decoder_seq_limit is None and self.moba_decoder_top_k_left is not None: + self.moba_use_decoder_seq_limit = self.moba_decoder_top_k_left * self.moba_block_size + self.check_legality_parameters() + + def check_legality_parameters( + self, + ) -> None: + if self.moba_encoder_top_k_left is not None: + assert self.moba_encoder_top_k_left > 0, "moba_encoder_top_k_left must large than 0" + + if self.moba_encoder_top_k_right is not None: + assert self.moba_encoder_top_k_right > 0, "moba_encoder_top_k_right must large than 0" + assert ( + self.moba_encoder_top_k_right >= self.moba_encoder_top_k_left + ), "moba_encoder_top_k_right must large than moba_encoder_top_k_left" + + if self.moba_decoder_top_k_left is not None: + assert self.moba_decoder_top_k_left > 0, "moba_decoder_top_k_left must large than 0" + + if self.moba_decoder_top_k_right is not None: + assert self.moba_decoder_top_k_right > 0, "moba_decoder_top_k_right must large than 0" + assert ( + self.moba_decoder_top_k_right >= self.moba_decoder_top_k_left + ), "moba_decoder_top_k_right must large than moba_decoder_top_k_left" + + if self.moba_use_encoder_seq_limit is not None and self.moba_encoder_top_k_left is not None: + assert self.moba_use_encoder_seq_limit >= self.moba_encoder_top_k_left * self.moba_block_size + if self.moba_use_decoder_seq_limit is not None and self.moba_decoder_top_k_left is not None: + assert self.moba_use_decoder_seq_limit >= self.moba_decoder_top_k_left * self.moba_block_size + + def to_json_string(self): + """ + Convert moba_attention_config to json string. + """ + return json.dumps({key: value for key, value in self.__dict__.items() if value is not None}) + + class EarlyStopConfig: def __init__( self, @@ -1038,6 +1099,7 @@ def __init__( decoding_config: DecodingConfig = None, quant_config: QuantConfigBase = None, graph_opt_config: GraphOptimizationConfig = None, + moba_attention_config: MobaAttentionConfig = None, speculative_config: SpeculativeConfig = None, tokenizer: str = None, max_model_len: int = 8192, @@ -1072,7 +1134,7 @@ def __init__( self.early_stop_config: Optional[EarlyStopConfig] = early_stop_config self.decoding_config: DecodingConfig = decoding_config # type: ignore self.cache_config: CacheConfig = cache_config # type: ignore - + self.moba_attention_config: Optional[MobaAttentionConfig] = moba_attention_config # Initialize cuda graph capture list if self.graph_opt_config.cudagraph_capture_sizes is None: self.graph_opt_config._set_cudagraph_sizes(max_num_seqs=self.parallel_config.max_num_seqs) @@ -1171,6 +1233,15 @@ def postprocess(self): self.paddle_commit_id = paddle.version.commit + if self.cache_config.enable_chunked_prefill: + self.force_chunked_prefill = int(envs.FD_FORCE_CHUNKED_PREFILL) + if ( + self.speculative_config is not None + and self.speculative_config.method in ["mtp"] + and not self.force_chunked_prefill + ): + self.cache_config.enable_chunked_prefill = False + if self.max_num_batched_tokens is None: if self.cache_config.enable_chunked_prefill: self.max_num_batched_tokens = 2048 From 8694219c696d40817deb8b45e4aef9a8177c6a7f Mon Sep 17 00:00:00 2001 From: ltd0924 <32387785+ltd0924@users.noreply.github.com> Date: Mon, 1 Sep 2025 23:23:15 +0800 Subject: [PATCH 5/8] Update fused_moe_backend_base.py --- .../model_executor/layers/moe/fused_moe_backend_base.py | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/fastdeploy/model_executor/layers/moe/fused_moe_backend_base.py b/fastdeploy/model_executor/layers/moe/fused_moe_backend_base.py index 82df14a36b4..73c9b634a8f 100644 --- a/fastdeploy/model_executor/layers/moe/fused_moe_backend_base.py +++ b/fastdeploy/model_executor/layers/moe/fused_moe_backend_base.py @@ -58,7 +58,7 @@ def init_ep(self, layer: nn.Layer) -> None: layer.ep_size, layer.ep_rank, layer.fd_config.model_config.redundant_experts_num, - layer.fd_config.parallel_config.ep_group, + ep_group=layer.fd_config.parallel_config.ep_group, ) self.ep_decoder_runner = EPDecoderRunner( layer.top_k, @@ -69,7 +69,7 @@ def init_ep(self, layer: nn.Layer) -> None: layer.ep_size, layer.ep_rank, layer.fd_config.model_config.redundant_experts_num, - layer.fd_config.parallel_config.ep_group, + ep_group=layer.fd_config.parallel_config.ep_group, ) else: if layer.fd_config.parallel_config.moe_phase.phase == "prefill": @@ -84,7 +84,7 @@ def init_ep(self, layer: nn.Layer) -> None: layer.ep_size, layer.ep_rank, layer.fd_config.model_config.redundant_experts_num, - layer.fd_config.parallel_config.ep_group, + ep_group=layer.fd_config.parallel_config.ep_group, ) else: from .ep import EPDecoderRunner @@ -98,7 +98,7 @@ def init_ep(self, layer: nn.Layer) -> None: layer.ep_size, layer.ep_rank, layer.fd_config.model_config.redundant_experts_num, - layer.fd_config.parallel_config.ep_group, + ep_group=layer.fd_config.parallel_config.ep_group, ) def process_loaded_weights(self, layer, weights) -> None: From 4e26c1c85a088bc5db69eef96fe14bbdf27c230a Mon Sep 17 00:00:00 2001 From: ltd0924 <32387785+ltd0924@users.noreply.github.com> Date: Mon, 1 Sep 2025 23:26:49 +0800 Subject: [PATCH 6/8] Update worker_process.py --- fastdeploy/worker/worker_process.py | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/fastdeploy/worker/worker_process.py b/fastdeploy/worker/worker_process.py index 3ee3e2647d1..f42eae5d066 100644 --- a/fastdeploy/worker/worker_process.py +++ b/fastdeploy/worker/worker_process.py @@ -290,7 +290,7 @@ def event_loop_normal(self) -> None: if self.fd_config.load_config.dynamic_load_weight: if self.parallel_config.enable_expert_parallel: - paddle.distributed.barrier(group=self.parallel_config.ep_group) + paddle.distributed.barrier(self.parallel_config.ep_group) else: paddle.distributed.barrier(self.parallel_config.tp_group) if self.model_weights_signal[0] != 0: @@ -307,7 +307,6 @@ def event_loop_normal(self) -> None: self.parallel_config.engine_worker_queue_port, ) self.model_weights_signal[0] = 0 - paddle.distributed.broadcast(self.model_weights_signal, src=0) if self.exist_task_signal.value[0] == 1 or self.task_queue.read_finish_flag.get() == 1: logger.info(f"Rank: {self.local_rank} Detected new requests.") From c8c1ecc9392b95156763c2dec936a521094dff32 Mon Sep 17 00:00:00 2001 From: ltd0924 <32387785+ltd0924@users.noreply.github.com> Date: Mon, 1 Sep 2025 23:36:36 +0800 Subject: [PATCH 7/8] Update worker_process.py --- fastdeploy/worker/worker_process.py | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/fastdeploy/worker/worker_process.py b/fastdeploy/worker/worker_process.py index f42eae5d066..b4d227e4e91 100644 --- a/fastdeploy/worker/worker_process.py +++ b/fastdeploy/worker/worker_process.py @@ -267,10 +267,6 @@ def event_loop_normal(self) -> None: if self.fd_config.load_config.dynamic_load_weight: paddle.distributed.broadcast(self.model_weights_signal, src=0, group=self.parallel_config.tp_group) - if self.parallel_config.tensor_parallel_size > 1: - # Synchronize before updating weights - paddle.distributed.barrier(self.parallel_config.tp_group) - self.insert_step = False req_dicts = None local_rank = self.local_rank % self.parallel_config.tensor_parallel_size @@ -288,6 +284,10 @@ def event_loop_normal(self) -> None: else: self.exist_task_signal.value[0] = 1 + if self.parallel_config.tensor_parallel_size > 1: + # Synchronize the signal for other workers + paddle.distributed.barrier(self.parallel_config.tp_group) + if self.fd_config.load_config.dynamic_load_weight: if self.parallel_config.enable_expert_parallel: paddle.distributed.barrier(self.parallel_config.ep_group) From 4d3d62efcdde036cd8fac73c626caf37445ad4ba Mon Sep 17 00:00:00 2001 From: ltd0924 <32387785+ltd0924@users.noreply.github.com> Date: Tue, 2 Sep 2025 14:42:38 +0800 Subject: [PATCH 8/8] Update dynamic_weight_manager.py --- fastdeploy/rl/dynamic_weight_manager.py | 1 + 1 file changed, 1 insertion(+) diff --git a/fastdeploy/rl/dynamic_weight_manager.py b/fastdeploy/rl/dynamic_weight_manager.py index d00db929140..66136a94af4 100644 --- a/fastdeploy/rl/dynamic_weight_manager.py +++ b/fastdeploy/rl/dynamic_weight_manager.py @@ -117,6 +117,7 @@ def clear_parameters(self, pid: int = 0) -> None: paddle.distributed.shutdown_process_group(self.parallel_config.tp_group) if self.parallel_config.enable_expert_parallel: paddle.distributed.barrier(self.parallel_config.ep_group) + paddle.distributed.shutdown_process_group(self.parallel_config.ep_group) self._update_shared_status(pid, -2) def _update_model_from_state(self, state_dict: Dict[str, paddle.Tensor], src_type: str):