pyzmq设置SNDHWM高水位标记后PUSH socket未按预期阻塞问题
ZMQ PUSH Socket设置HWM未按预期阻塞的问题解决
问题概述
- 预期行为:PUSH Socket发送1500条消息后阻塞,等待PULL Socket消费完消息再继续发送
- 实际行为:PUSH Socket发送约7万条消息后才出现阻塞
- 运行后关键输出:
... Message [72747::1678881796.1216621] Message [72748::1678881796.1216683] Message [72749::1678881796.1216772] Message [72750::1678881796.121687] Message [72751::1678881796.1216931] Socket is blocked! Waiting for the worker to consume messages... Worker is slow! Waiting... Worker is slow! Waiting...
原始代码
PUSH端代码
import time import zmq context = zmq.Context() # 设置发送高水位标记为1500的PUSH Socket socket = context.socket(zmq.PUSH) socket.setsockopt(zmq.SNDHWM, 1500) socket.bind("tcp://127.0.0.1:5566") # 发送消息循环 for i in range(200000): try: print(f"Message [{i}::{time.time()}]") socket.send_string(f"Message [{i}::{time.time()}]", zmq.DONTWAIT) except zmq.error.Again: # 处理HWM触发的阻塞情况 print("Socket is blocked! Waiting for the worker to consume messages...") while True: # 轮询检查Socket是否可发送 if socket.poll(timeout=5000, flags=zmq.POLLOUT): break else: print("Worker is slow! Waiting...") time.sleep(1) # 资源清理 socket.close() context.term()
PULL端代码
import time import zmq context = zmq.Context() receiver = context.socket(zmq.PULL) receiver.setsockopt(zmq.RCVHWM, 1) receiver.connect("tcp://127.0.0.1:5566") def worker(): i = 0 while True: try: message = receiver.recv_string(zmq.NOBLOCK) print(f"Received message: [{i}:{time.time()}]{message}") time.sleep(500e-3) i+=1 except zmq.Again: print("no messages") time.sleep(100e-3) worker()
问题原因
- TCP缓冲区缓存溢出:ZMQ发送消息时会先写入操作系统的TCP发送缓冲区,默认缓冲区容量较大,能容纳大量小消息。即使ZMQ的SNDHWM达到1500上限,只要TCP缓冲区未满,
send_string(zmq.DONTWAIT)就不会抛出zmq.Again异常。 - 无连接时的消息堆积:PUSH端先绑定端口后立刻发送消息,此时PULL端可能尚未建立连接,ZMQ会将消息暂存在内部队列中。当PULL连接建立后,ZMQ会批量将内部队列的消息发送到TCP缓冲区,进一步扩大了实际发送量。
解决方案
1. 限制TCP发送缓冲区大小
在PUSH端设置TCP发送缓冲区的容量,避免系统层缓存过多消息:
# 按单条消息大小计算,确保缓冲区仅能容纳约1500条消息 single_msg_size = len("Message [0::1678881796.1216621]") socket.setsockopt(zmq.SNDBUF, single_msg_size * 1500) # 或直接设置固定值,比如50KB socket.setsockopt(zmq.SNDBUF, 1024 * 50)
2. 等待PULL连接建立后再发送
修改PUSH端代码,确保PULL已连接再开始发送,避免无连接时的消息堆积:
# 绑定后等待PULL连接 print("Waiting for PULL connection...") while True: if socket.poll(timeout=100, flags=zmq.POLLOUT): break
3. 启用IMMEDIATE选项
开启ZMQ_IMMEDIATE,让PUSH仅在有活跃连接时才发送消息,减少内部队列的无意义堆积:
socket.setsockopt(zmq.IMMEDIATE, 1)
修改后的完整代码
PUSH端代码
import time import zmq context = zmq.Context() socket = context.socket(zmq.PUSH) # 设置发送高水位标记 socket.setsockopt(zmq.SNDHWM, 1500) # 限制TCP发送缓冲区大小 socket.setsockopt(zmq.SNDBUF, 1024 * 50) # 启用IMMEDIATE,仅在有连接时发送消息 socket.setsockopt(zmq.IMMEDIATE, 1) socket.bind("tcp://127.0.0.1:5566") # 等待PULL端建立连接 print("Waiting for PULL connection...") while True: if socket.poll(timeout=100, flags=zmq.POLLOUT): break # 发送消息 for i in range(200000): try: print(f"Message [{i}::{time.time()}]") socket.send_string(f"Message [{i}::{time.time()}]", zmq.DONTWAIT) except zmq.error.Again: print("Socket is blocked! Waiting for the worker to consume messages...") while True: if socket.poll(timeout=5000, flags=zmq.POLLOUT): break else: print("Worker is slow! Waiting...") time.sleep(1) socket.close() context.term()
PULL端代码(优化为阻塞接收)
import time import zmq context = zmq.Context() receiver = context.socket(zmq.PULL) receiver.setsockopt(zmq.RCVHWM, 1) receiver.connect("tcp://127.0.0.1:5566") def worker(): i = 0 while True: # 使用阻塞接收替代非阻塞,减少空轮询消耗 message = receiver.recv_string() print(f"Received message: [{i}:{time.time()}]{message}") time.sleep(500e-3) i += 1 worker()
内容的提问来源于stack exchange,提问作者LordTitiKaka
相关产品推荐
相关产品推荐

