惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

Apple Machine Learning Research
Apple Machine Learning Research
M
MIT News - Artificial intelligence
罗磊的独立博客
博客园 - 【当耐特】
A
About on SuperTechFans
Last Week in AI
Last Week in AI
雷峰网
雷峰网
IT之家
IT之家
aimingoo的专栏
aimingoo的专栏
H
Hackread – Cybersecurity News, Data Breaches, AI and More
博客园_首页
博客园 - 叶小钗
Microsoft Azure Blog
Microsoft Azure Blog
博客园 - Franky
J
Java Code Geeks
Cyber Security Advisories - MS-ISAC
Cyber Security Advisories - MS-ISAC
D
Docker
Engineering at Meta
Engineering at Meta
B
Blog RSS Feed
The Cloudflare Blog
大猫的无限游戏
大猫的无限游戏
阮一峰的网络日志
阮一峰的网络日志
S
SegmentFault 最新的问题
Recent Announcements
Recent Announcements

又见苍岚

COLMAP PatchMatch Stereo 算法详解 事件驱动的状态机框架:从理论到工程实践 Git 在国内网络环境下无法 Push 的排查与修复 —— 配置 Clash 代理 分段五次多项式插值原理详解 路径插值方法深度对比研究 Claude Code 使用指南 OpenClaw 记忆管理与技能创建指南 CBS(Conflict-Based Search)算法详解 A* 算法及其变种详解 OpenClaw 配置多 Agents Windows Powershell 无法加载文件,因为在此系统上禁止运行脚本问题的解决方案 MaxClaw 安装流程 大模型 AI 名词介绍 AList 网盘聚合工具简介 Protobuf 简介与测试 Claude Code 简介以及 GLM 4.7 模型接入 Github 歌词下载工具 163MusicLyrics Python __getattr__ 懒加载 Python TypedDict 机器人仿真平台 Gazebo 安装记录 机器人仿真平台 Gazebo 简介 多机器人路径规划问题(Multi-Agent Path Finding, MAPF)简介 Python exifread 读取修改过的 jpeg 信息错误问题修复 3D 坐标系变换的理解 3D 旋转矩阵基本概念 MongoDB Compass 介绍 Python 环境管理工具 uv Flutter 开发指南 Snipaste 安装下载与黑屏问题解决方案 全局路径规划算法记录
ZMQ 代理 XPUB/XSUB 模式
Yiwei Zhang · 2025-03-05 · via 又见苍岚

本文详细记录 ZeroMQ 的代理 XPUB/XHUB 模式。

ZeroMQ 的 XPUB/XSUB 模式是构建消息代理系统的核心模式,常用于构建发布-订阅架构的中介代理。

核心概念说明

  1. XPUB (Extended PUB)

    • 特殊类型的发布者套接字
    • 可以接收订阅者的订阅/取消订阅请求
    • 会以二进制消息形式通知上游发布者订阅变化
  2. XSUB (Extended SUB)

    • 特殊类型的订阅者套接字
    • 可以显式发送订阅/取消订阅请求
    • 支持动态修改订阅过滤条件
  3. 典型拓扑结构

Python 代码示例

1. 代理服务器(Broker)

1
2
3
4
5
6
7
8
9
10
11
12
13
14
import zmq

context = zmq.Context()

# 前端接收发布者消息(XSUB)
frontend = context.socket(zmq.XSUB)
frontend.bind("tcp://*:5555") # 发布者连接到此端口

# 后端发送给订阅者(XPUB)
backend = context.socket(zmq.XPUB)
backend.bind("tcp://*:5556") # 订阅者连接到此端口

# 使用代理进行消息转发
zmq.proxy(frontend, backend)

2. 发布者(Publisher)

1
2
3
4
5
6
7
8
9
10
11
import zmq
import time

context = zmq.Context()
socket = context.socket(zmq.PUB)
socket.connect("tcp://localhost:5555") # 连接到代理的XSUB端口

while True:
socket.send_string("A/当前时间 %s" % time.ctime())
socket.send_string("B/当前时间 %s" % time.ctime())
time.sleep(1)

3. 订阅者(Subscriber)

1
2
3
4
5
6
7
8
9
10
11
12
import zmq

context = zmq.Context()
socket = context.socket(zmq.SUB)
socket.connect("tcp://localhost:5556") # 连接到代理的XPUB端口

# 订阅以"A/"开头的消息
socket.setsockopt_string(zmq.SUBSCRIBE, "A/")

while True:
message = socket.recv_string()
print("收到消息:", message)

tcp://*:5555

  • 含义
    * 表示绑定到本机的所有可用网络接口(即所有 IP 地址)。例如:
    • 如果你的机器有 IP 192.168.1.10010.0.0.2,绑定到 * 后,外部可以通过这两个 IP 访问服务。
    • 如果替换为具体 IP(如 tcp://192.168.1.100:5555),则只允许通过该 IP 访问。
  • 典型场景
    当需要让其他机器连接到你的服务时,必须使用 *。如果仅在本地测试,可以用 tcp://127.0.0.1:5555(仅允许本机连接)。

关键特性说明

  1. 订阅通知机制

    • XPUB 会收到二进制格式的订阅通知
    • 首字节 \x01 表示订阅,\x00 表示取消订阅
    • 例如:b'\x01A/' 表示订阅以 “A/” 开头的消息
  2. 动态订阅管理

    1
    2
    3
    # 订阅者可以动态修改订阅条件
    socket.setsockopt_string(zmq.UNSUBSCRIBE, "A/")
    socket.setsockopt_string(zmq.SUBSCRIBE, "B/")
  3. 多对多通信

    • 允许多个发布者和多个订阅者同时连接
    • 自动处理消息路由
  4. bindconnect 的区别

    操作 角色 行为 典型用途
    bind 服务端/监听方 创建一个监听端口,等待他人连接 代理(Broker)、服务提供者
    connect 客户端/连接方 主动连接到一个已存在的地址 发布者(Publisher)、订阅者(Sub)

    关键差异:

    1. 顺序要求
      • 必须先有 bind(服务端),然后才能 connect(客户端)。
      • 如果客户端尝试连接未绑定的地址,会失败(Connection refused)。
    2. 灵活性
      • bind 方是稳定的服务节点(如消息代理)。
      • connect 方是动态的客户端节点(可以随时加入或离开)。
    3. 地址所有权
      • bind 的地址是独占的(同一端口只能被一个进程绑定)。
      • connect 可以多对一(多个客户端连接同一个服务端)。实际应用场景
  5. 消息总线(Message Bus)

  6. 服务解耦中间层

  7. 分布式日志系统

  8. 实时数据广播系统

常见问题处理

  1. 慢订阅者问题

    • 使用 zmq.CONFLATE=1 选项保留最新消息
    • 设置 HWM(高水位线)控制队列长度
  2. 调试技巧

    1
    2
    # 监控XPUB的订阅请求
    backend.setsockopt(zmq.XPUB_VERBOSE, 1)
  3. 性能优化

    • 使用 inproc 传输进行进程内通信
    • 启用多线程代理

文章链接:
https://www.zywvvd.com/notes/tools/zmq/zmq-proxy/zmq-proxy/