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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 09:38:59