You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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()

问题原因

  1. TCP缓冲区缓存溢出:ZMQ发送消息时会先写入操作系统的TCP发送缓冲区,默认缓冲区容量较大,能容纳大量小消息。即使ZMQ的SNDHWM达到1500上限,只要TCP缓冲区未满,send_string(zmq.DONTWAIT)就不会抛出zmq.Again异常。
  2. 无连接时的消息堆积: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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.28 09:45:07