[Bug Fix] fix bug for PD EP (#4823)

* fix bug for PD EP

* fix

* optimize perf for engine worker queue

* fix bug

* fix internode ll two stage

* fix for ci

* fix bug
This commit is contained in:
chenjian
2025-11-10 15:33:29 +08:00
committed by GitHub
parent 112623e33e
commit 78895e2c7d
11 changed files with 109 additions and 44 deletions

View File

@@ -556,10 +556,13 @@ class PaddleDisWorkerProc:
def start_task_queue_service(self):
# Initialize task queue
task_address = (
self.parallel_config.pod_ip,
self.parallel_config.engine_worker_queue_port,
)
if not envs.FD_ENGINE_TASK_QUEUE_WITH_SHM:
task_address = (
self.parallel_config.pod_ip,
self.parallel_config.engine_worker_queue_port,
)
else:
task_address = f"/dev/shm/fd_task_queue_{self.parallel_config.engine_worker_queue_port}.sock"
logger.info(f"connect task queue address {task_address}")
self.task_queue = TaskQueue(
address=task_address,