理解 ZMQ 的 HWM

Gui*_*imo 6 python zeromq pyzmq

我在理解 ZeroMQ 高水位线 (HWM) 队列的工作原理时遇到了一些麻烦。

我在下面附上了两个脚本,它们重现了以下内容。

  • 建立 PUSH/PULL 连接,将所有 HWM 队列设置为大小 1。
  • 让拉马者睡一会儿。
  • 从推送器发送 2200 条消息。
  • 当 puller 醒来时,接收 2200 条消息并打印它们。

我得到的结果是 puller 能够成功接收(打印)所有消息。此外,推送器似乎几乎立即完成执行。根据ZMQ 官方文档,我期望推送程序不会在 puller 唤醒之前完成执行,因为send(...)由于到达 HWM 而在第二次调用时被阻塞。我还尝试在每次send(...)通话之间添加 0.001 秒的睡眠,结果相同。

所以,我的问题是:

  • 在send(...)达到 HWM(大小 1)后,为什么 pusher 在第二次调用中没有阻塞?
  • pusher 和 puller 中存储的消息在哪里?
  • HWM 大小与存储的消息数量之间是否存在直接关系?

脚本:

推手.py

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)

拉手.py

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。

Jos*_*and 9

这里至少涉及 4 个缓冲区:zmq 发送缓冲区、OS 写入 tcp 缓冲区、OS 读取 tcp 缓冲区和 zmq 接收缓冲区。

当消息成功写入操作系统 tcp 写入缓冲区时,zmq io 线程会将消息标记为“已发送”。这些消息现在被视为“在传输中”。

然后网络堆栈负责将尽可能多的传输到另一个进程的匹配 OS recv 缓冲区,最后,接收 zmq io 线程一次最多从该缓冲区读取 HWM 消息到 ZMQ recv 队列。

默认情况下,操作系统缓冲区通常在 10-100kb 左右,并且在 ZMQ 甚至注意到对方没有消耗任何消息之前,这两个缓冲区都可以完全填满“传输中”消息。出于性能原因,这些缓冲区是必需的 - 您不能只是摆脱它们。

您的问题的解决方案可能涉及 req/rep 套接字和明确的应用程序级 ack,即指南中的惰性海盗模式。