【问题标题】:0mq one-to-many connection0mq 一对多连接
【发布时间】:2011-11-02 22:49:26
【问题描述】:

使用0mq在进程之间建立双向通信最正确的方法是什么?我需要创建几个后台进程来等待来自主进程的命令,执行一些计算并将结果返回给主进程。

【问题讨论】:

    标签: python multithreading zeromq


    【解决方案1】:

    有几种方法可以做到这一点。最直接的方法可能是使用REQ/REP 套接字。每个后台进程/worker 都会有一个 REP 套接字,您可以使用 REQ 套接字与它们通信:

    import zmq
    
    def worker(addr):
        context = zmq.Context()
        socket = context.socket(zmq.REP)
        socket.bind(addr)
        while True:
            # get message from boss
            msg = socket.recv()
            # ...do smth
            # send back results
            socket.send(msg)
    
    if __name__ == '__main__':
        # spawn 5 workers
        from multiprocessing import Process
        for i in range(5):
            Process(target=worker, args=('tcp://127.0.0.1:500%d' % i,)).start()
    

    您必须连接到每个工作人员才能向他们发送消息并返回结果:

    context = zmq.Context()
    socket = context.socket(zmq.REQ)
    socket.connect(worker_addr)
    socket.send('message')
    msg = socket.recv()
    

    另一种方法是使用PUB/SUB 向工作人员发送消息并使用PUSH/PULL 来获取结果:

    import zmq
    
    def worker(worker_id, publisher_addr, results_addr):
        context = zmq.Context()
        sub = context.socket(zmq.SUB)
        sub.connect(publisher_addr)
        sub.setsockopt(zmq.SUBSCRIBE, worker_id)
        push = context.socket(zmq.PUSH)
        push.connect(results_addr)
    
        while True:
            msg = sub.recv_multipart()[1]
            # do smth, send off results
            push.send_multipart([worker_id, msg])
    
    if __name__ == '__main__':
        publisher_addr = 'tcp://127.0.0.1:5000'
        results_addr = 'tcp://127.0.0.1:5001'
    
        # launch some workers into space
        from multiprocessing import Process
        for i in range(5):
            Process(target=worker, args=('worker-%d' % i, publisher_addr, results_addr,)).start()
    

    要将命令广播给特定的工作人员,您可以执行以下操作:

    context = zmq.Context()
    pub = context.socket(zmq.PUB)
    pub.bind(publisher_addr)
    # send message to worker-1
    pub.send_multipart(['worker-1', 'hello'])
    

    拉入结果:

    context = zmq.Context()
    pull = context.socket(zmq.PULL)
    pull.bind(results_addr)
    
    while True:
        worker_id, result = pull.recv_multipart()
        print worker_id, result
    

    【讨论】:

      【解决方案2】:

      考虑使用Request Reply Broker,但将 REQ 套接字交换为 DEALER。 DEALER 不会阻塞发送,并会自动对您的工作人员的流量进行负载平衡。

      图片中的Client 是你的main processService A/B/C 是你的background processes (workers)Main process 应该绑定到一个端点。 Workers 应该连接到主进程的端点以接收工作项。

      main process 中保留工作项列表和发送时间。如果一段时间没有答复,请重新发送工作项,因为 worker 可能已经死了。

      【讨论】:

        猜你喜欢
        • 2023-03-16
        • 2013-05-17
        • 1970-01-01
        • 2016-08-11
        • 2020-05-25
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多