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

Asyncio环境下ZMQ PUB/SUB订阅端无法接收消息问题排查

Asyncio环境下ZMQ PUB/SUB订阅端无法接收消息的解决方案

问题场景

在asyncio(搭配uvloop)环境中搭建ZMQ发布-订阅(PUB/SUB)架构,订阅端基于asyncio实现,但无法接收发送端发送的任何消息。

订阅端代码片段

async def zmq_listener(sub, stop_event):
    global mode
    while not stop_event.is_set():
        msg = await sub.recv_multipart()
        print(msg)
        if len(msg) == 1:
            print(msg[0].decode())

######################################
### MAIN PROGRAM STARTS HERE 
async def main():
    tasks = []
    tasks.append(asyncio.create_task(other_routine1(abc, stop_event)))

    if doZMQ:
        print("ZMQ Mode is enabled")
        ctx = zmq.asyncio.Context()
        sub = ctx.socket(zmq.SUB)
        ip = 'tcp://127.0.0.1:5559'
        sub.connect(ip)
        sub.setsockopt(zmq.SUBSCRIBE, b"") # Subscribe to all topics        
        tasks.append(asyncio.create_task(zmq_listener(sub, stop_event)))

    await asyncio.gather(*tasks)

stop_event = asyncio.Event()

try:
    uvloop.install()
    asyncio.run(main(), debug=False)

except KeyboardInterrupt:
    print("Program interrupted by user")
    asyncio.run(stop())

测试发送端代码

import zmq

ctx = zmq.Context()
pub = ctx.socket(zmq.PUB)
ip = 'tcp://127.0.0.1:5559'
pub.bind(ip)

pub.send_multipart([b"SHOW"])
pub.send_multipart([b"P1"])

核心原因

  1. 发送端消息发送时机过早:ZMQ的PUB套接字不会为未完成订阅握手的客户端缓存消息。发送端启动后立即发送消息时,订阅端可能还在建立TCP连接、发送订阅指令的过程中,导致消息直接被丢弃。
  2. 资源清理不规范:订阅端未正确关闭ZMQ套接字和上下文,虽不直接导致收不到消息,但可能引发后续稳定性问题。

修复方案

1. 调整发送端,等待订阅端完成连接

在发送端绑定端口后添加延迟,确保订阅端完成连接和订阅握手:

import zmq
import time

ctx = zmq.Context()
pub = ctx.socket(zmq.PUB)
ip = 'tcp://127.0.0.1:5559'
pub.bind(ip)

# 等待订阅端完成连接握手
time.sleep(1)

pub.send_multipart([b"SHOW"])
pub.send_multipart([b"P1"])

# 等待消息发送完成再清理资源
time.sleep(0.1)
pub.close()
ctx.term()

2. 优化订阅端的资源清理逻辑

在main函数中添加finally块,确保程序退出时正确关闭ZMQ资源:

async def main():
    tasks = []
    tasks.append(asyncio.create_task(other_routine1(abc, stop_event)))
    sub = None
    ctx = None

    if doZMQ:
        print("ZMQ Mode is enabled")
        ctx = zmq.asyncio.Context()
        sub = ctx.socket(zmq.SUB)
        ip = 'tcp://127.0.0.1:5559'
        sub.connect(ip)
        sub.setsockopt(zmq.SUBSCRIBE, b"") # Subscribe to all topics        
        tasks.append(asyncio.create_task(zmq_listener(sub, stop_event)))

    try:
        await asyncio.gather(*tasks)
    finally:
        # 清理ZMQ资源
        if sub:
            sub.close()
        if ctx:
            ctx.term()

3. 严格执行启动顺序

必须先启动订阅端,待其完成连接和订阅后,再启动发送端发送消息,确保消息能被正常接收。

关键原理说明

ZMQ的PUB/SUB模式中,PUB套接字仅向已完成订阅握手的SUB套接字转发消息。订阅端需要完成TCP连接建立、发送订阅指令两个步骤后,才能成为PUB的有效订阅者,此时发送的消息才会被转发。

内容的提问来源于stack exchange,提问作者HermDP

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 02:15:36