【问题标题】:ZeroMQ inter process communication from single thread loses messages来自单线程的 ZeroMQ 进程间通信丢失消息
【发布时间】:2016-04-20 22:04:49
【问题描述】:

我目前正在探索测试我的 zeromq 应用程序的可能性。我的印象是我可以在同一个线程中有一个发布者/订阅者,让发布者发布和订阅者订阅它而不会丢失消息。然而,当我让发布者发送几条消息时,没有一条消息能传递给订阅者。

这是我使用的代码:

import zmq

def main():
    ctx = zmq.Context.instance()
    sender = ctx.socket(zmq.PUB)
    sender.setsockopt(zmq.HWM, 1000)
    sender.bind('tcp://*:10001')

    rcvr = ctx.socket(zmq.SUB)
    rcvr.setsockopt(zmq.HWM, 1000)
    rcvr.connect('tcp://127.0.0.1:10001')
    rcvr.setsockopt(zmq.SUBSCRIBE, "")

    for i in range(100):
        sender.send('%i' % i)

    while True:
        try:
            print rcvr.recv(zmq.NOBLOCK)
        except zmq.ZMQError:
            break


if __name__ == '__main__':
    main()

运行此程序时,我没有得到任何输出。

让我印象深刻的是,接收方在发送方发送之前就已连接,因此应该对这些消息进行排队。或者这是一个完全错误的假设,我应该使用 PUSH/PULL 代替?

【问题讨论】:

  • 检查guide并搜索slow joiner。

标签: python unit-testing zeromq pyzmq


【解决方案1】:

我认为这是 ZeroMQ guide 中描述的慢连接器问题。

这种“慢速加入者”症状经常出现在足够多的人身上,我们将对其进行详细解释。

我认为主要问题是在订阅者套接字开始侦听之前所有消息都已发送,并且消息飞过并被丢弃。在设置套接字和发送消息之间设置延迟不起作用,因为在接收器开始侦听之前已经发送了最后一条消息。

正如您所建议的,推/拉套接字在内存中执行队列作业。您可以像这样在单个进程中的套接字之间发送作业

# pushpull.py
import zmq

def main():
    ctx = zmq.Context()
    sender = ctx.socket(zmq.PUSH)
    sender.bind('tcp://*:10001')

    rcvr = ctx.socket(zmq.PULL)
    rcvr.connect('tcp://127.0.0.1:10001')

    for i in range(100):
        sender.send_unicode('%i' % i)

    while True:
        msg = rcvr.recv()
        print(msg)

if __name__ == '__main__':
    main()

或者,如果您想使用 pub/sub 套接字,我们需要两个进程和一个 time.sleep(1) 在套接字设置和消息发送之间:

首先启动接收器

# rcvr.py
import zmq

def main():
    ctx = zmq.Context()
    rcvr = ctx.socket(zmq.SUB)
    rcvr.connect('tcp://127.0.0.1:10001')
    rcvr.setsockopt_string(zmq.SUBSCRIBE, "")

    while True:
        msg = rcvr.recv()
        print(msg)

if __name__ == '__main__':
    main()

然后是发件人,

# sender.py
import zmq
import time

def main():
    ctx = zmq.Context()
    sender = ctx.socket(zmq.PUB)
    sender.bind('tcp://*:10001')

    time.sleep(1)
    for i in range(100):
        sender.send_unicode('%i' % i)

if __name__ == "__main__":
    main()

接收:

b'0'
b'1'
b'2'
b'3' ...

我目前正在 Python 3.3 和 pyzmq 13.1.0 中使用出色的 WinPython 分发版,因此 zmq 调用中的一些字符串处理以及打印功能有些不同。 希望对您有所帮助。

【讨论】:

    【解决方案2】:

    您应该将 SUB 套接字连接到端口 10000 而不是 10001。目前 SUB 套接字正在等待发布者,而 PUB 套接字正在等待订阅者。 0mq 允许“客户端”在不存在“服务器”的情况下进行连接的功能也意味着当您连接到端口 10001 时不会引发错误,这是设计使然。

    【讨论】:

    • 很好地发现了它!然而,这是一个错字,并不能解决问题。 :(
    【解决方案3】:

    让我印象深刻的是,接收者在发送者发送之前已经连接

    这实际上不是真的 - 接收者已经开始了连接过程,但这并不意味着该过程已经完成。连接是异步的。

    如果您实际上将其用于进程内通信,我建议您使用inproc 传输,这不是问题:

    url = 'inproc://whatever'
    sender.bind(url)
    ...
    recvr.connect(url)
    

    【讨论】:

    • 然而,它没有理由不能与 TCP 一起工作,不是吗?
    • 我所说的其他内容,非进程连接是异步的,因此您在实际连接的对等方之前开始发送,因此消息被丢弃。如果您在发送前有一个睡眠,它会正常工作。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2011-07-03
    • 1970-01-01
    • 2011-11-20
    • 1970-01-01
    • 2011-02-09
    • 2021-07-24
    • 1970-01-01
    相关资源
    最近更新 更多