Gui*_*imo 6 python zeromq pyzmq
我在理解 ZeroMQ 高水位线 (HWM) 队列的工作原理时遇到了一些麻烦。
我在下面附上了两个脚本,它们重现了以下内容。
我得到的结果是 puller 能够成功接收(打印)所有消息。此外,推送器似乎几乎立即完成执行。根据ZMQ 官方文档,我期望推送程序不会在 puller 唤醒之前完成执行,因为send(...)由于到达 HWM 而在第二次调用时被阻塞。我还尝试在每次send(...)通话之间添加 0.001 秒的睡眠,结果相同。
所以,我的问题是:
send(...)达到 HWM(大小 1)后,为什么 pusher 在第二次调用中没有阻塞?脚本:
import zmq
context = zmq.Context()
push_socket = context.socket(zmq.PUSH)
push_socket.setsockopt(zmq.SNDHWM, 1)
push_socket.setsockopt(zmq.RCVHWM, 1)
push_socket.bind("tcp://127.0.0.1:5557")
print(push_socket.get_hwm()) # Prints 1
print('Sending all messages')
for i in range(2200):
push_socket.send(str(i).encode('ascii'))
print('Finished execution...')
Run Code Online (Sandbox Code Playgroud)
import zmq
import time
context = zmq.Context()
pull_socket = context.socket(zmq.PULL)
pull_socket.setsockopt(zmq.RCVHWM, 1)
pull_socket.setsockopt(zmq.SNDHWM, 1)
pull_socket.connect("tcp://127.0.0.1:5557")
print(pull_socket.get_hwm()) # Prints 1
print('Connected, but not receiving yet... (Sleep 4s)')
time.sleep(4)
print('Receiving everything now!')
rec = ''
for i in range(2200):
rec += '{} '.format(pull_socket.recv().decode('ascii'))
print(rec) # Prints `0 1 2 ... 2198 2199 `
Run Code Online (Sandbox Code Playgroud)
为了重现我的测试用例,打开两个终端并在一个终端中启动第一个 puller.py,然后在另一个终端中快速启动(4 秒窗口)pusher.py。
这里至少涉及 4 个缓冲区:zmq 发送缓冲区、OS 写入 tcp 缓冲区、OS 读取 tcp 缓冲区和 zmq 接收缓冲区。
当消息成功写入操作系统 tcp 写入缓冲区时,zmq io 线程会将消息标记为“已发送”。这些消息现在被视为“在传输中”。
然后网络堆栈负责将尽可能多的传输到另一个进程的匹配 OS recv 缓冲区,最后,接收 zmq io 线程一次最多从该缓冲区读取 HWM 消息到 ZMQ recv 队列。
默认情况下,操作系统缓冲区通常在 10-100kb 左右,并且在 ZMQ 甚至注意到对方没有消耗任何消息之前,这两个缓冲区都可以完全填满“传输中”消息。出于性能原因,这些缓冲区是必需的 - 您不能只是摆脱它们。
您的问题的解决方案可能涉及 req/rep 套接字和明确的应用程序级 ack,即指南中的惰性海盗模式。