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

ZeroMQ发布端重启后订阅端失效问题求助

ZeroMQ订阅端在PUB重启后丢失订阅的解决方案

问题根源

ZeroMQ的PUB-SUB模式是无状态的,当PUB端(如你的监视器)重启时,SUB端(诊断端)已建立的连接会失效,但ZeroMQ默认不会主动检测连接状态、重新建立连接并恢复订阅,这就是需要重启诊断端才能恢复的原因。


解决方案1:启用ZeroMQ自动重连机制

ZeroMQ内置了自动重连功能,通过设置套接字参数让SUB在连接断开后自动尝试重连,重连成功后会自动重新发送订阅指令。

修改诊断端的初始化代码,添加重连参数:

def zeromq_init():
    # 统一使用一个Context,推荐单进程单Context
    ZMQdata.context = zmq.Context()
    
    # 订阅水泵
    ZMQdata.in_socket1 = ZMQdata.context.socket(zmq.SUB)
    subscribe_to = f"tcp://{PUMP_ADDR}:{PUMP_TO_MONITOR_PORT}"
    # 设置自动重连参数
    ZMQdata.in_socket1.setsockopt(zmq.RECONNECT_IVL, 1000)  # 初始重连间隔1秒
    ZMQdata.in_socket1.setsockopt(zmq.RECONNECT_IVL_MAX, 5000)  # 最大重连间隔5秒
    ZMQdata.in_socket1.connect(subscribe_to)
    ZMQdata.in_socket1.setsockopt_string(zmq.SUBSCRIBE, "")
    
    # 订阅监视器
    ZMQdata.in_socket2 = ZMQdata.context.socket(zmq.SUB)
    subscribe_to = f"tcp://{MONITOR_ADDR}:{MONITOR_TO_PUMP_PORT}"
    ZMQdata.in_socket2.setsockopt(zmq.RECONNECT_IVL, 1000)
    ZMQdata.in_socket2.setsockopt(zmq.RECONNECT_IVL_MAX, 5000)
    ZMQdata.in_socket2.connect(subscribe_to)
    ZMQdata.in_socket2.setsockopt_string(zmq.SUBSCRIBE, "")
  • 无需手动处理重连逻辑,ZeroMQ会自动维护连接,监视器重启后SUB会自动恢复连接和订阅。

解决方案2:结合心跳超时主动重连

利用你已有的心跳机制,在诊断端监听心跳超时,主动重建套接字和订阅,比纯依赖自动重连更贴合业务场景。

示例代码(需结合你的心跳处理逻辑):

import time
import zmq

# 初始化时记录最后收到心跳的时间
ZMQdata.last_monitor_heartbeat = time.time()
ZMQdata.last_pump_heartbeat = time.time()

def check_connections():
    # 检查监视器连接
    try:
        # 非阻塞接收消息
        msg = ZMQdata.in_socket2.recv(zmq.NOBLOCK)
        # 处理心跳消息,更新时间
        if msg == b"MONITOR_HEARTBEAT":
            ZMQdata.last_monitor_heartbeat = time.time()
        # 处理其他消息...
    except zmq.Again:
        # 30秒未收到心跳,判定连接失效
        if time.time() - ZMQdata.last_monitor_heartbeat > 30:
            # 关闭旧套接字,重新初始化
            ZMQdata.in_socket2.close()
            ZMQdata.in_socket2 = ZMQdata.context.socket(zmq.SUB)
            ZMQdata.in_socket2.setsockopt(zmq.RECONNECT_IVL, 1000)
            ZMQdata.in_socket2.setsockopt(zmq.RECONNECT_IVL_MAX, 5000)
            ZMQdata.in_socket2.connect(f"tcp://{MONITOR_ADDR}:{MONITOR_TO_PUMP_PORT}")
            ZMQdata.in_socket2.setsockopt_string(zmq.SUBSCRIBE, "")
            ZMQdata.last_monitor_heartbeat = time.time()
    
    # 水泵连接的检查逻辑同上...

# 在主循环中定期调用check_connections()
while True:
    check_connections()
    time.sleep(1)

解决方案3:引入XPUB-XSUB代理(最健壮架构)

通过中间代理层解耦PUB和SUB,所有PUB连接到代理的XSUB端口,所有SUB连接到代理的XPUB端口。代理会自动维护所有连接,SUB无需关心后端PUB的状态变化。

代理端代码(可运行在任意一台树莓派上):

import zmq

context = zmq.Context()
# XSUB接收所有PUB的消息
xsub_socket = context.socket(zmq.XSUB)
xsub_socket.bind("tcp://*:5556")
# XPUB转发消息给所有SUB
xpub_socket = context.socket(zmq.XPUB)
xpub_socket.bind("tcp://*:5557")

# 启动代理
zmq.proxy(xsub_socket, xpub_socket)

修改水泵/监视器的PUB端代码:

将原来绑定端口改为连接到代理的XSUB端口:

# 水泵的PUB示例
pub_socket = zmq.Context().socket(zmq.PUB)
pub_socket.connect("tcp://{PROXY_ADDR}:5556")  # PROXY_ADDR是代理所在树莓派的IP

修改诊断端的SUB代码:

连接到代理的XPUB端口,只需连接一次,无需关心后端PUB的重启:

def zeromq_init():
    ZMQdata.context = zmq.Context()
    # 订阅代理的XPUB端口,同时接收水泵和监视器的消息
    ZMQdata.in_socket = ZMQdata.context.socket(zmq.SUB)
    ZMQdata.in_socket.connect(f"tcp://{PROXY_ADDR}:5557")
    ZMQdata.in_socket.setsockopt_string(zmq.SUBSCRIBE, "")
  • 优势:彻底解决PUB重启导致的订阅丢失问题,架构扩展性强,后续新增PUB/SUB只需连接代理即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 22:16:20