Fix zmq binding (#2930)

Co-authored-by: Chunyuan WU <chunyuan.wu@intel.com>
This commit is contained in:
Lianmin Zheng
2025-01-16 14:36:07 -08:00
committed by GitHub
co-authored by Chunyuan WU
parent bf3edc2c60
commit 0427416b59
5 changed files with 18 additions and 12 deletions
@@ -66,7 +66,7 @@ class DataParallelController:
self.context = zmq.Context(1 + server_args.dp_size)
if server_args.node_rank == 0:
self.recv_from_tokenizer = get_zmq_socket(
self.context, zmq.PULL, port_args.scheduler_input_ipc_name
self.context, zmq.PULL, port_args.scheduler_input_ipc_name, False
)
# Dispatch method
@@ -93,6 +93,7 @@ class DataParallelController:
self.context,
zmq.PUSH,
dp_port_args[dp_rank].scheduler_input_ipc_name,
True,
)
def launch_dp_schedulers(self, server_args, port_args):
@@ -58,10 +58,10 @@ class DetokenizerManager:
# Init inter-process communication
context = zmq.Context(2)
self.recv_from_scheduler = get_zmq_socket(
context, zmq.PULL, port_args.detokenizer_ipc_name
context, zmq.PULL, port_args.detokenizer_ipc_name, True
)
self.send_to_tokenizer = get_zmq_socket(
context, zmq.PUSH, port_args.tokenizer_ipc_name
context, zmq.PUSH, port_args.tokenizer_ipc_name, False
)
if server_args.skip_tokenizer_init:
+4 -4
View File
@@ -162,21 +162,21 @@ class Scheduler:
if self.attn_tp_rank == 0:
self.recv_from_tokenizer = get_zmq_socket(
context, zmq.PULL, port_args.scheduler_input_ipc_name
context, zmq.PULL, port_args.scheduler_input_ipc_name, False
)
self.send_to_tokenizer = get_zmq_socket(
context, zmq.PUSH, port_args.tokenizer_ipc_name
context, zmq.PUSH, port_args.tokenizer_ipc_name, False
)
if server_args.skip_tokenizer_init:
# Directly send to the TokenizerManager
self.send_to_detokenizer = get_zmq_socket(
context, zmq.PUSH, port_args.tokenizer_ipc_name
context, zmq.PUSH, port_args.tokenizer_ipc_name, False
)
else:
# Send to the DetokenizerManager
self.send_to_detokenizer = get_zmq_socket(
context, zmq.PUSH, port_args.detokenizer_ipc_name
context, zmq.PUSH, port_args.detokenizer_ipc_name, False
)
else:
self.recv_from_tokenizer = None
@@ -119,10 +119,10 @@ class TokenizerManager:
# Init inter-process communication
context = zmq.asyncio.Context(2)
self.recv_from_detokenizer = get_zmq_socket(
context, zmq.PULL, port_args.tokenizer_ipc_name
context, zmq.PULL, port_args.tokenizer_ipc_name, True
)
self.send_to_scheduler = get_zmq_socket(
context, zmq.PUSH, port_args.scheduler_input_ipc_name
context, zmq.PUSH, port_args.scheduler_input_ipc_name, True
)
# Read model args