ØMQ.rs
纯 Rust ZeroMQ:用于分布式和并发应用的无代理消息传递。在进程内、进程间以及网络上表现一致的套接字级消息传递模式。
- 适用于 Linux、macOS 和 Windows 的 Tokio 后端
- 20 种套接字类型:稳定的 ZMQ 模式以及草案 CLIENT/SERVER、RADIO/DISH、SCATTER/GATHER、CHANNEL/PEER 和 STREAM
- 8 种稳定传输协议:TCP、IPC、inproc、UDP、WS、WSS、
lz4+tcp://、 以及lz4+ws://;实验性的zstd+tcp://压缩传输协议位于zstd特性之后 - 3 种安全机制:NULL、PLAIN、CURVE
- 无需 C 编译器,无需 libzmq,无需 libsodium
- Python 绑定(pyomq),C API(omq-libzmq)
REQ/REP 延迟
PUB/SUB 吞吐量
LZ4 PUSH/PULL 吞吐量
难点
OMQ 旨在实现真实的 ZMQ 行为,而不仅仅是理想路径下的 PUSH/PULL 吞吐量。您将获得:
- 无需额外调优的 ZeroMQ 语义:没有拓扑特定的 socket 类型,没有用户可见的批处理 API,没有手动重连循环。
- 传输故障是常态:重连、先连接后绑定、对端变动以及绑定端重启都是设计的一部分。
- 对端故障不会变成用户错误:
send()和recv()在断开连接、重连、慢消费者以及绑定端重启期间保持正常工作。 - 在负载下的 HWM 背压和路由公平性,而不仅仅是在空队列示例中。
- 关于无对端发送、linger 和 HWM 的已记录 libzmq 兼容性边缘情况: doc/libzmq/semantics.md.
- 热路径具有大小感知和低延迟意识:小消息保持内联且无分配,inproc 按值传递消息,大负载在关键处使用零拷贝缓冲区。
- 唯一遵循 libzmq 架构的 Rust ZeroMQ 实现:应用线程与专用的后台 IO 线程分离,IO 工作在这些线程上线性扩展,PUB 对端自动分配给 IO 通道。
- 公共 crate 使用内存安全的 Rust。
unsafe被隔离并使用 Miri 进行检查。 - 基准测试覆盖真实形态:CPU 核算、扇入/扇出、公平性、传输差异。
用法
[!NOTE] API 仍在演变,可能在次要版本之间发生变化。欢迎提交 bug 报告并在真实工作负载中进行测试。
Rust 后端是 omq-tokio:在 Linux、
macOS 和 Windows 上使用 tokio + mio。它可以在 OMQ 拥有的运行时线程上
或现有的 tokio 运行时内运行套接字 IO。
| API | 运行时位置 | 扩展模型 |
|---|---|---|
Context::new().socket(...) | 异步套接字,OMQ 拥有的 IO 线程 | 线性 IO 通道扩展;PUB 对等节点跨通道分片 |
Context::current().socket(...) | 异步套接字,调用者的活动 tokio 运行时 | Tokio 调度器/工作窃取;PUB 扇出保留在一个 OMQ 通道上 |
Context::new().blocking_socket(...) | 同步套接字,OMQ 拥有的 IO 线程 | 线性 IO 通道扩展;调用者线程不介入 IO |
支持的 32 位 Linux 目标是 i686-unknown-linux-gnu 和
armv7-unknown-linux-gnueabihf。它们需要原生 64 位原子操作。ZMTP 线
长度字段保持 64 位,但实际的帧/消息大小受
平台分配限制约束(在 32 位上低于 4 GiB)。
如果你了解 ZeroMQ,你就了解 OMQ。相同的套接字类型,相同的 connect/bind/send/recv:
use omq_tokio::{Context, Message, Options, SocketType};
let ctx = Context::new();
let push = ctx.socket(SocketType::Push, Options::default());
push.connect("tcp://127.0.0.1:5555".parse()?).await?;
push.send(Message::single("hello")).await?;
let pull = ctx.socket(SocketType::Pull, Options::default());
pull.bind("tcp://127.0.0.1:5555".parse()?).await?;
let msg = pull.recv().await?;
assert_eq!(&msg[0], b"hello");
更多示例见 examples/zguide-tokio/, 这是 ZeroMQ Guide 模式到 OMQ 的移植。
Cargo 特性
均为可选。默认构建为最小部署:NULL 机制 + TCP / IPC / inproc / UDP,无需 C 编译器。可启用以下任意一项:
| feature | 添加内容 | 额外依赖 |
|---|---|---|
plain | PLAIN 用户名/密码认证(RFC 24) | - |
curve | CURVE 加密握手机制(RFC 26) | crypto_box, crypto_secretbox |
lz4 | lz4+tcp:// 压缩传输(RFC) | lz4rip |
zstd | 实验性 zstd+tcp:// 压缩传输 | zrip |
ws | WebSocket(ws://)和 安全 WebSocket(wss://)传输 | rustls, rustls-native-certs |
设计亮点
| 特性 | 详情 |
|---|---|
Sans-I/O ZMTP 编解码器 (omq-proto) | 字节输入 / 事件输出;热路径上无异步,无 trait。镜像 rustls::ConnectionCommon。 |
| 消息计数 HWM | send_hwm/recv_hwm 计数完整消息,而非字节。发送 HWM 针对每个出站管道/环形缓冲区,因此当存在多个管道或发送槽时,原生缓冲的总消息数可能超过一个 send_hwm。 |
| 连续帧负载 | &msg[0] 直接提供 &[u8];无可能失败的借用,无合并步骤。 |
| 零拷贝发送和接收 | 发送:大型 Bytes 负载无需任何数据拷贝即可到达内核 writev。接收:大型帧直接读取到预分配的缓冲区中,绕过中间队列。 |
| Patricia 树订阅匹配器 | 复杂度为 O(M)(M 为主题长度),而非 O(NxM)。 |
| LZ4 字典自动训练 | 默认关闭。启用后,从前 100 条消息中进行训练,并向对端发送一次;将有效压缩阈值从 512 B 降低至 64 B。 |
| 监控事件 | 类套接字的 Stream,在每个连接 / 断开连接 / 握手事件上拥有 PeerInfo。 |
工作区
五个 Cargo 工作区 crate 加上 Python 绑定。
| Crate | 功能 | 不安全策略 |
|---|---|---|
omq-proto | 无 I/O 的 ZMTP 3.x 核心:编解码器、消息、机制、订阅 | #![forbid(unsafe_code)] |
omq-tokio | 多线程 tokio 后端(Linux/macOS/Windows) | #![forbid(unsafe_code)] |
omq-libzmq | 兼容 libzmq 的 C 接口(libomq_zmq 动态/静态库) | 不安全的 C ABI 边界 |
yring | 有界 SPSC 环形缓冲区,采用 ypipe 风格的批量刷新/预取 | 不安全的环形核心,经 Miri 测试 |
omq-bench | 基准测试运行器和 SVG 图表生成器 | 仅限基准测试的进程控制和 CPU 核算 |
pyomq | Python 绑定(基于 omq-tokio 的 PyO3,同步 + asyncio) | PyO3 FFI 边界 |
测试
每种套接字类型、传输、机制和功能组合均 由集成测试覆盖。测试套件是分层的:
- 700+ Rust tests 涵盖 socket 类型、传输层、机制以及 libzmq 兼容的 C API 行为。
- Feature-gated coverage 针对 PLAIN、CURVE、LZ4 以及 pyzmq/libzmq 互操作性。WebSocket 拥有专门的测试和 soak 覆盖。
- Protocol fuzzing(默认 opt-in 运行约 100 万次迭代, 可配置更长的运行时间):针对 wire 解析器和 socket-action 状态机的手工编写模糊测试。
- 20+ soak scenarios 涵盖 Rust 和 pyomq:peer churn、重连 风暴、PUB/SUB churn、ROUTER/DEALER churn、HWM 重连、取消 安全性、压缩(lz4)、PLAIN / CURVE 认证、机制重连、 大消息吞吐量、多 socket、inproc 跨线程、 WebSocket 吞吐量和重连。Soak 运行会采样 RSS 和 FD 计数。
- Loom 覆盖
yringSPSC 内存排序、异步唤醒以及omq-tokio信号竞争窗口。 - Miri 针对
yring。 - Release semver review 通过
release-plz。
./scripts/test-all.sh # standard sweep with local perf gate
OMQ_FUZZ=1 ./scripts/test-all.sh # include fuzz suites
OMQ_SKIP_PYOMQ=1 ./scripts/test-all.sh
OMQ_SKIP_PERF=1 ./scripts/test-all.sh
浸泡测试有意与完整扫描分离:
FEATURES="soak lz4 plain curve ws"
OMQ_SOAK_DURATION_SECS=600 cargo test -p omq-tokio \
--features "$FEATURES" --release --test omq_soak_peer_churn -- --nocapture
延伸阅读
- COMPARISONS.md:跨实现对比图表。
- BENCHMARKS_COMPRESSION.md:在带宽受限链路上的 lz4+tcp 吞吐量。
- doc/architecture.md:架构及 tokio 后端内部机制。
- doc/libzmq/semantics.md:关于无对端发送、linger 和 HWM 的精确兼容性 说明。
- doc/lz4-rfc.md:LZ4 压缩传输线路 格式及字典分发规则。
平台与要求
Linux 是主要的开发和基准测试平台。 PR CI 必需检查涵盖 Linux x86_64、macOS ARM64 和 Windows Rust 测试, 以及 32 位 Linux 交叉检查。Fmt 和 clippy 也在 macOS Intel 上运行。 扩展 CI 涵盖 Ubuntu ARM64。macOS 作业串行运行 Rust 测试, 因为在托管运行器上 socket/定时器时序更为敏感。
macOS 在 CI 中同时覆盖 Intel 和 ARM64 运行器。
omq-tokio 使用 mio / kqueue。omq-libzmq 使用基于管道的
通知 fd 用于 zmq_poll/ZMQ_FD 就绪状态。
Windows 在 CI 中覆盖。omq-tokio 支持 TCP、IPC 命名
管道、inproc、UDP 和 WebSocket 传输。omq-libzmq 在 Windows 上构建并
测试受支持的 C API 表面。
pyomq 发布 Linux、macOS 和 Windows 的 wheels 以及 sdist。
要求:
- Rust 1.93 或更高版本(edition 2024)。
贡献
请参阅 CONTRIBUTING.md 获取指南,以及 DEVELOPMENT.md 获取构建、测试和基准测试命令。
AI 披露
本项目在架构、实现、测试、基准测试基础设施和文档方面均获得了大量 LLM 辅助。它是一次关于 LLM 辅助开发能做什么和不能做什么的实验。设计决策和方向由我负责。
License
ISC.