【问题标题】:Get data from ZMQ PULL socket. How to sync computation?从 ZMQ PULL 套接字获取数据。如何同步计算?
【发布时间】: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


    【解决方案1】:

    您正在寻求实施障碍。

    ZeroMQ 完全是关于 Actor 模型编程的,其特点之一是在发送和接收消息时没有显式隐含的会合。也就是说,无论另一端是否已阅读消息,发送都会返回。

    所以这意味着必须在 ZeroMQ 的 Actor 模型之上合成一个屏障(一种会合类型)。

    1. 使用 PUSH / PULL 套接字对将数据发送给工作人员。
    2. 为工作人员使用单独的 PUSH/PULL 套接字对向生产者发回“我有数据并准备好继续”消息。
    3. 让生产者等待这些“我可以继续”消息,
    4. 当它从每个工作人员那里收到一个消息时,在 PUB / SUB 套接字上向工作人员发送一条“go”消息。

    通信顺序进程

    出于兴趣,您可能希望将 Actor 模型编程与通信顺序进程(在 Rust、Erlang 和(我认为?)Go 中正在卷土重来)进行比较。在 CSP 中发送/接收消息是一个会合。这有几个好处;

    • 发件人知道已收到消息,而不仅仅是排队,
    • 如果有性能和延迟目标,它会强制用户正确解决架构和资源分配问题。您无法隐藏传输中的消息。因此,如果没有提供足够的工人,生产者显然无法卸载消息;增加的延迟不能暂时隐藏缺陷。
    • 如果您设法构建了一个可以死锁、活锁等的架构,它总是会的。然而,Actor 模型架构多年来似乎都非常好,直到有一天网络变得有点繁忙。

    要使用 CSP 做您想做的事,您可以省略上面的第 2 步和第 3 步。 Producer 会在发送到最后一个 worker 返回时知道每个 worker 都收到了它的数据,并且可以立即发送“go”。

    就个人而言,我真的希望 ZeroMQ 可以选择成为 CSP,而不是 Actor。那么这将是美妙的,而不是非常巨大。它真正好的原因在于它是否是 tcp、ipc、inproc 等都无关紧要。它的行为都是相同的(显然速度变化)。

    AFAIK Rust、Erlang 和 Go CSP 通道仅此过程。 ZMQ 可以是进程间和/或进程内和/或计算机间,这使得它非常适合开发可能超出一台计算机的系统。需要将线程卸载到另一台计算机?更改连接字符串,无需更改其他代码。很不错。

    【讨论】:

      【解决方案2】:

      您正在为命令和数据使用单独的流 - 这将始终保证同步问题。在接收方,您将有两个流缓冲区 - 第一个有大量数据要处理,第二个只有命令和 poll() 将确保通知您两者都已准备好读取。

      我看到了两种处理这个问题的方法:

      1) 保持简单:只使用一个流。您最后发送的所有内容最终都会收到。 TCP 保证了这一点。如果你使用的是 json,你可以添加 'type': 'command' 或 'type': 'data' 来区分消息类型。

      2) 如果出于某种原因,您确实需要两个流(例如,您确实想使用发布者/订阅者模式),则接收者应在发送者发送其命令之前确认发送者接收到最后一批数据。如果 所有 工作人员需要在 任何 工作人员使用该命令启动之前接收他们的数据,则也可以选择此选项。

      【讨论】:

      • 嗯,1)我看到了一个问题。我有多个 PULLER(我想扇出数据),如何确保将数据与推送一起发送给所有消费者(我认为不可能)?我想我需要使用解决方案 2,但是使用循环,我不知道某些数据是否会最后发送给消费者。
      猜你喜欢
      • 1970-01-01
      • 2016-12-17
      • 2017-06-29
      • 2014-03-05
      • 1970-01-01
      • 2023-03-23
      • 2018-12-23
      • 2012-09-14
      • 2014-10-29
      相关资源
      最近更新 更多