连接XPUB-XSUB代理时,发布端如何等待订阅转发?
问题描述
我用ZeroMQ的XPUB-XSUB代理构建消息总线,支持订阅者和发布者动态增减,用于多个守护进程通信。但遇到一个问题:当已有订阅者连接到代理的XPUB端时,新发布者连接后立即发消息,第一条消息会丢失。
我推测原因是发布者连接时,订阅者信息还没及时同步到发布端,导致第一条消息被套接字丢弃。目前用了一个不可靠的临时方法:在发布端连接后加短暂sleep再发消息。
想问有没有可靠的方式等待订阅转发完成?或者是否应该换用其他类型的套接字?
示例代码
import threading import time from zmq import Context, Socket, proxy from zmq.constants import PUB, SUB, XPUB, XPUB_VERBOSE, XSUB def message_bus(): context = Context.instance() in_socket: Socket = context.socket(XSUB) in_socket.bind("ipc:///tmp/in_socket.ipc") out_socket: Socket = context.socket(XPUB) out_socket.bind("ipc:///tmp/out_socket.ipc") out_socket.setsockopt(XPUB_VERBOSE, True) proxy(in_socket, out_socket) def publisher(): context = Context.instance() bus_in_socket: Socket = context.socket(PUB) bus_in_socket.connect("ipc:///tmp/in_socket.ipc") count = 1 while True: bus_in_socket.send_string(f"message number {count}") count += 1 time.sleep(0.5) def subscriber(): context = Context.instance() bus_out_socket: Socket = context.socket(SUB) bus_out_socket.connect("ipc:///tmp/out_socket.ipc") bus_out_socket.subscribe("") while True: print(f"subscriber {bus_out_socket.recv_multipart()}") if __name__ == "__main__": message_bus_thread = threading.Thread(target=message_bus, daemon=True) subscriber_thread = threading.Thread(target=subscriber, daemon=True) publisher_thread = threading.Thread(target=publisher, daemon=True) message_bus_thread.start() subscriber_thread.start() time.sleep(1) publisher_thread.start() message_bus_thread.join() subscriber_thread.join() publisher_thread.join()
输出结果
subscriber [b'message number 2'] subscriber [b'message number 3'] subscriber [b'message number 4'] subscriber [b'message number 5'] subscriber [b'message number 6'] subscriber [b'message number 7'] subscriber ...
可见第一条消息[b'message number 1']完全未被接收。
解决方案
1. 利用XPUB订阅通知实现同步
XPUB套接字在开启XPUB_VERBOSE后,会向代理发送所有订阅/取消订阅事件消息。你可以替换默认的proxy()函数,自定义代理逻辑,实现订阅状态的同步通知:
- 用自定义消息循环处理XSUB和XPUB之间的消息流转
- 维护订阅者计数,当XPUB收到订阅消息(以
0x01开头的帧),除了转发到XSUB,还可以向发布者发送订阅就绪的信号 - 发布者连接后,先等待代理的订阅就绪信号,再发送第一条业务消息
2. 发布者主动监听订阅事件
发布者可以额外创建一个SUB套接字,连接到代理的XPUB端,监听订阅事件,确认有订阅者存在后再发送消息:
def publisher(): context = Context.instance() # 业务发布套接字 bus_in_socket: Socket = context.socket(PUB) bus_in_socket.connect("ipc:///tmp/in_socket.ipc") # 监听订阅事件的套接字 sub_listener = context.socket(SUB) sub_listener.connect("ipc:///tmp/out_socket.ipc") sub_listener.subscribe("") # 等待订阅事件(XPUB发送的订阅帧以0x01开头) while True: msg = sub_listener.recv_multipart() if msg[0].startswith(b'\x01'): break count = 1 while True: bus_in_socket.send_string(f"message number {count}") count += 1 time.sleep(0.5)
3. 换用带确认机制的套接字类型
如果需要绝对可靠的消息投递,PUB/SUB的广播模式本身不保证消息接收,可以考虑以下替代方案:
- ROUTER/DEALER:支持请求-响应模式,发布者发送消息后可等待订阅者的确认,确保消息被接收
- Push/Pull:适合任务分发场景,可自行添加确认逻辑增强可靠性
- 若仍需广播能力,PUB/SUB+代理的模式依然适用,但必须处理订阅同步问题
4. 等待套接字连接就绪(比sleep更可靠)
用poll()替代固定sleep,等待发布套接字连接建立完成:
bus_in_socket.connect("ipc:///tmp/in_socket.ipc") # 等待套接字可写(连接就绪) poller = zmq.Poller() poller.register(bus_in_socket, zmq.POLLOUT) poller.poll(1000) # 超时1秒
注意:此方法仅确保连接建立,无法保证订阅信息同步,仍可能丢失第一条消息,但比固定sleep更灵活。
内容的提问来源于stack exchange,提问作者Afkaaja
相关产品推荐
相关产品推荐

