【发布时间】:2017-02-08 05:25:17
【问题描述】:
我有一个生产者使用 PULL / PUSH 向多个工人发送数据。在执行计算任务之前,所有工作人员都需要接收所有数据。
我尝试使用发送“go”的 PUB / SUB 套接字进行同步,但由于 PUSH 套接字是非阻塞的,因此在数据流结束之前接收到 go...
发件人:
context = zmq.Context()
push_socket = self.context.socket(zmq.PUSH)
push_socket.bind("tcp://127.0.0.1:5557")
pull_socket = self.context.socket(zmq.PULL)
pull_socket.bind("tcp://127.0.0.1:5558")
for index, data in range(100):
push_socket.send_json({"data": data, "id": index})
pub_socket.send_json({"command": "map"})
接收者:
# recieve work
consumer_receiver = context.socket(zmq.PULL)
consumer_receiver.connect("tcp://127.0.0.1:5557")
# receive commands
consumer_command = context.socket(zmq.SUB)
consumer_command.subscribe("")
consumer_command.connect("tcp://127.0.0.1:5559")
poller = zmq.Poller()
poller.register(consumer_receiver, zmq.POLLIN)
poller.register(consumer_command, zmq.POLLIN)
while True:
events = dict(poller.poll(100))
if consumer_command in events:
received = consumer_command.recv_json()
command = received["command"]
print("received command : ", command)
if consumer_receiver in events:
received = consumer_receiver.recv_json()
print("received data", received)
接收器输出:
received data {'data': ['Hi'], 'id': 0}
received command : map
received data {'data': ['hi'], 'id': 1}
...
我想要:
received data {'data': ['Hi'], 'id': 0}
received data {'data': ['hi'], 'id': 1}
...
received command : map
我尝试将 PUSH 套接字的 HWM 设置为 1,但没有成功。
如何在 PUSH 完成后向所有工作人员发送同步消息?
【问题讨论】:
标签: python sockets network-programming zeromq pyzmq