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

推荐订阅源

Hugging Face - Blog
Hugging Face - Blog
Google DeepMind News
Google DeepMind News
云风的 BLOG
云风的 BLOG
WordPress大学
WordPress大学
Vercel News
Vercel News
Apple Machine Learning Research
Apple Machine Learning Research
T
Tailwind CSS Blog
I
InfoQ
小众软件
小众软件
Recent Announcements
Recent Announcements
博客园 - 【当耐特】
The GitHub Blog
The GitHub Blog
大猫的无限游戏
大猫的无限游戏
美团技术团队
T
The Blog of Author Tim Ferriss
奇客Solidot–传递最新科技情报
奇客Solidot–传递最新科技情报
酷 壳 – CoolShell
酷 壳 – CoolShell
MongoDB | Blog
MongoDB | Blog
V
V2EX
J
Java Code Geeks
有赞技术团队
有赞技术团队
博客园 - 聂微东
B
Blog RSS Feed
博客园 - 司徒正美

又见苍岚

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/