如何在ZMQ/0MQ中构建多发布者多订阅者网络?是否需消息代理?
嘿,针对你问的ZMQ/0MQ多发布者多订阅者架构的问题,我来一步步给你拆解,结合你提供的代码示例来讲解~
ZMQ多发布者+多订阅者架构实现指南
一、核心问题解答:是否必须用消息代理?
不是必须的!ZMQ提供了两种核心方案:直接连接的无代理模式,以及基于XPUB-XSUB代理的灵活模式,你可以根据场景选择。
二、无代理的直接连接架构(适合发布者地址固定场景)
如果你的发布者地址/端口是固定可访问的,订阅者可以直接连接多个发布者,无需中间代理。
完善后的代码示例
我把你提供的代码补全并扩展,实现两个发布者、两个订阅者的直接连接:
import time import zmq from multiprocessing import Process def bind_pub(sleep_seconds, max_messages, pub_id, port): context = zmq.Context() socket = context.socket(zmq.PUB) # 每个发布者绑定独立端口 socket.bind(f"tcp://*:{port}") message = 0 while message < max_messages: # 用前缀做主题,方便订阅者过滤特定发布者的消息 socket.send_string(f"pub_{pub_id} message_number={message}") print(f"Publisher {pub_id} sent message {message}") message += 1 time.sleep(sleep_seconds) socket.close() context.term() def sub_process(sub_id, *pub_ports): context = zmq.Context() socket = context.socket(zmq.SUB) # 订阅者同时连接多个发布者的端口 for port in pub_ports: socket.connect(f"tcp://localhost:{port}") # 订阅所有消息(也可以指定主题,比如只订阅pub_1的消息:socket.setsockopt_string(zmq.SUBSCRIBE, "pub_1")) socket.setsockopt_string(zmq.SUBSCRIBE, "") message_count = 0 while message_count < 10: # 接收10条消息后退出 msg = socket.recv_string() print(f"Subscriber {sub_id} received: {msg}") message_count += 1 socket.close() context.term() if __name__ == "__main__": # 启动两个发布者,分别绑定5556、5557端口 pub1 = Process(target=bind_pub, args=(1, 5, 1, 5556)) pub2 = Process(target=bind_pub, args=(1, 5, 2, 5557)) # 启动两个订阅者,同时连接两个发布者 sub1 = Process(target=sub_process, args=(1, 5556, 5557)) sub2 = Process(target=sub_process, args=(2, 5556, 5557)) pub1.start() pub2.start() sub1.start() sub2.start() pub1.join() pub2.join() sub1.join() sub2.join()
关键说明
- 每个发布者绑定独立端口,避免端口冲突;
- 利用主题前缀(比如
pub_1)实现订阅者对特定发布者消息的过滤; - 缺点:如果发布者动态增减(比如频繁上线下线),订阅者需要手动维护连接列表,扩展性较差。
三、基于XPUB-XSUB代理的灵活架构(适合动态场景)
当发布者是动态的,或者不想让订阅者感知所有发布者的地址,就可以用ZMQ的XPUB-XSUB代理作为中间枢纽,统一处理消息转发。
代码示例
import time import zmq from multiprocessing import Process def proxy_process(): context = zmq.Context() # XSUB端口:接收所有发布者的消息 xsub_socket = context.socket(zmq.XSUB) xsub_socket.bind("tcp://*:5558") # XPUB端口:向所有订阅者转发消息 xpub_socket = context.socket(zmq.XPUB) xpub_socket.bind("tcp://*:5559") # 启动代理,自动完成消息转发 zmq.proxy(xsub_socket, xpub_socket) xsub_socket.close() xpub_socket.close() context.term() def pub_via_proxy(sleep_seconds, max_messages, pub_id): context = zmq.Context() socket = context.socket(zmq.PUB) # 发布者只需连接代理的XSUB端口 socket.connect("tcp://localhost:5558") message = 0 while message < max_messages: socket.send_string(f"pub_{pub_id} message_number={message}") print(f"Publisher {pub_id} (via proxy) sent message {message}") message += 1 time.sleep(sleep_seconds) socket.close() context.term() def sub_via_proxy(sub_id): context = zmq.Context() socket = context.socket(zmq.SUB) # 订阅者只需连接代理的XPUB端口 socket.connect("tcp://localhost:5559") socket.setsockopt_string(zmq.SUBSCRIBE, "") message_count = 0 while message_count < 10: msg = socket.recv_string() print(f"Subscriber {sub_id} (via proxy) received: {msg}") message_count += 1 socket.close() context.term() if __name__ == "__main__": # 先启动代理进程 proxy = Process(target=proxy_process) proxy.start() time.sleep(0.5) # 等待代理初始化完成 # 启动发布者和订阅者 pub1 = Process(target=pub_via_proxy, args=(1, 5, 1)) pub2 = Process(target=pub_via_proxy, args=(1, 5, 2)) sub1 = Process(target=sub_via_proxy, args=(1,)) sub2 = Process(target=sub_via_proxy, args=(2,)) pub1.start() pub2.start() sub1.start() sub2.start() pub1.join() pub2.join() sub1.join() sub2.join() # 终止代理进程 proxy.terminate() proxy.join()
关键说明
- 代理作为中间层,发布者和订阅者只需要和代理交互,无需知道彼此的存在;
- XPUB套接字会自动把订阅者的订阅信息转发给发布者,方便发布者做动态调整(比如只发送有订阅需求的主题);
- 优势:架构扩展性极强,支持发布者和订阅者的动态增减。
四、核心注意事项
- 主题过滤:PUB-SUB模式的核心特性,订阅者通过
SUBSCRIBE设置主题前缀来过滤消息,这是区分不同发布者消息的关键; - 慢订阅者问题:ZMQ的PUB套接字会丢弃无法及时发送的消息(无可靠投递保证),如果订阅者处理速度慢,可能会丢消息,必要时可以开启
ZMQ_CONFLATE选项或改用其他可靠模式; - 上下文管理:每个进程必须创建独立的
zmq.Context,不要跨进程共享Context对象,避免出现异常。
内容的提问来源于stack exchange,提问作者Greg
相关产品推荐
相关产品推荐

