【问题标题】:Does a ZeroMQ Router and Dealer have to be in the same (Python) process for two way communication to be achieved?ZeroMQ 路由器和经销商是否必须在同一个(Python)进程中才能实现双向通信?
【发布时间】:2015-06-15 04:15:30
【问题描述】:

这是一个简单的 req-rep 服务。通常 zmq.REP 和 zmq.REQ 就足够了,但这不是我正在构建的案例应用程序。下面是一个执行客户端服务器通信的脚本。服务器在一个线程中运行。通信按预期工作。客户端发送消息,服务器接收消息,然后发送消息给客户端,客户端接收消息。

import time
from threading import Thread
import zmq

def worker_thread():
    cxt = zmq.Context.instance()
    worker = cxt.socket(zmq.DEALER)
    worker.setsockopt(zmq.IDENTITY, 'A')
    worker.connect("tcp://127.0.0.1:5559")

    for _ in range(10):
        request = worker.recv()
        print 'worker recieved'
        worker.send_multipart(['A', "data_recieved"])

cxt = zmq.Context.instance()
client = cxt.socket(zmq.ROUTER)
client.bind('tcp://127.0.0.1:5559')

Thread(target=worker_thread).start()
time.sleep(2)

for _ in range(10):
    client.send_multipart(['A', 'data'])
    request = client.recv()
    print 'worker responded'

当这段代码被分成两个脚本,客户端和服务器时,通信失败。客户端发送消息,但服务器从不接收它。代码如下:

客户:

import time
import zmq

cxt = zmq.Context.instance()
client = cxt.socket(zmq.ROUTER)
client.bind('tcp://127.0.0.1:5559')

for _ in range(10):
    client.send_multipart(['A', 'data'])
    request = client.recv()
    print 'worker responded'

服务器:

import time
import zmq

cxt = zmq.Context.instance()
worker = cxt.socket(zmq.DEALER)
worker.setsockopt(zmq.IDENTITY, 'A')
worker.connect("tcp://127.0.0.1:5559")

for _ in range(10):
    request = worker.recv()
    print 'worker recieved'
    worker.send_multipart(['A', "data_recieved"])

我能想到的唯一可能导致通信失败的原因是 zmq.ROUTER 和 zmq.DEALER 在单独的脚本中运行,因此是单独的进程。

【问题讨论】:

  • 不,它们可以位于不同的进程、不同的可执行文件、不同的计算机或不同的网络上。其他一些问题

标签: python zeromq router


【解决方案1】:

我认为您只是错误地颠倒了这两种套接字类型。

Router 套接字的行为类似于 Rep 套接字,而 Dealer 的行为类似于 Req。详情见:http://zeromq.org/tutorials:dealer-and-router

Router 和 Dealer 的特殊之处在于它们允许异步消息传递,其中 Req/Rep 专门用于同步交换。我不知道为什么您的第一个脚本有效,但您可能应该尝试反转套接字。

这是一个 C# 中的示例,它在 python 中模仿您的示例。这是两个独立的过程。

流程A:

const string Endpoint = "tcp://127.0.0.1:5559";

void Main()
{
    using (var ctx = NetMQContext.Create())
    using (var worker = ctx.CreateRouterSocket())
    {
        worker.Connect(Endpoint);

        for (int i = 0; i < 10; i++)
        {
            NetMQMessage message = router.ReceiveMessage();
            Console.WriteLine("#{0}: worker received: {1}", i, string.Join(", ", message.Select(t => t.ConvertToString())));
            message.Clear();
            message.Append("A");
            message.Append("data_recieved");
            worker.SendMessage(message);
        }
    }
}

流程 B:

const string Endpoint = "tcp://127.0.0.1:5559";

void Main()
{
    using (var ctx = NetMQContext.Create())
    using (var client = ctx.CreateDealerSocket())
    {
        client.Options.Identity = Encoding.ASCII.GetBytes("A");
        client.Bind(Endpoint);

        var message = new NetMQMessage();
        for (int i = 0; i < 10; i++)
        {
            message.Clear();
            message.Append("A");
            message.Append("data");
            client.SendMessage(message);
            var response = client.ReceiveMessage();
            Console.WriteLine("#{0}: worker responded", i);
        }
    }
}

输出似乎是您所期望的:

流程A:

#0: worker received: A, A, data
#1: worker received: A, A, data
#2: worker received: A, A, data
#3: worker received: A, A, data
#4: worker received: A, A, data
#5: worker received: A, A, data
#6: worker received: A, A, data
#7: worker received: A, A, data
#8: worker received: A, A, data
#9: worker received: A, A, data

流程 B:

#0: worker responded
#1: worker responded
#2: worker responded
#3: worker responded
#4: worker responded
#5: worker responded
#6: worker responded
#7: worker responded
#8: worker responded
#9: worker responded

干杯

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-08-23
    • 1970-01-01
    • 1970-01-01
    • 2019-08-07
    • 2018-08-03
    • 1970-01-01
    相关资源
    最近更新 更多