高性能 Client/Server 消息中间件 — ZeroMQ ROUTER/DEALER + bcrypt 认证 + 实时监控
Project description
PulseMQ
面向金融行情的高性能消息中间件,基于 ZeroMQ ROUTER/DEALER 架构。
- Client/Server 架构 — 服务端(ROUTER)集中路由 + 控制面;客户端(DEALER)发布/订阅
- PLAIN + bcrypt 认证 — ZAP 认证链:bcrypt 哈希凭据存储,admin token 保护监控接口
- 类型保真 — DataFrame / dict / str / bytes 端到端保真(
_restore_type自动还原) - 完整监控 — Web UI(ECharts)+ REST API + SSE 实时推送,延迟分位、事件流、在线客户端
- 运行期重连 — 断线自动指数退避重连(1s → 2s → 4s ... → 30s 封顶)
安装
Python >= 3.13
pip install pulse-mq
依赖项:ZeroMQ、msgspec、python-snappy、lz4、zstandard、pyarrow、pandas、bcrypt、loguru。
import pulsemq # Python 模块名(无连字符)
from pulsemq import Server
PyPI 分发名是
pulse-mq(pip install用),import 名是pulsemq(因 Python 标识符不允许连字符)。与python-dateutil→import dateutil模式一致。
快速开始
启动服务端
from pulsemq import Server
srv = Server(
data_endpoint="tcp://0.0.0.0:5555",
control_endpoint="tcp://0.0.0.0:5556",
admin_endpoint="0.0.0.0:9090",
credentials={"user1": "pass1"}, # 或 credentials_file
)
await srv.start()
首次启动若不存在凭据文件,自动生成默认 admin 用户,密码输出到 stderr。
可通过 CLI 管理用户:
pulsemq users add user1 pass1 --role publisher
pulsemq users list
服务端内置定时推送
无需外部生产者客户端,直接在 Server 上注册定时回调:
from pulsemq import Server
srv = Server(data_endpoint=..., control_endpoint=..., credentials={"u": "p"})
@srv.producer("market.tick", interval=2.0, serializer="msgpack")
async def gen_tick():
return {"symbol": "AAPL", "price": 180.5, "volume": 1000}
@srv.producer("market.quote", interval=0.5, serializer="pyarrow", compression="lz4")
async def gen_quote():
import pandas as pd
return pd.DataFrame({"price": [10, 20], "vol": [100, 200]})
@srv.burst_producer("bench", serializer="msgpack")
async def bench():
if not has_more():
return None
return {"seq": next_seq()}
await srv.start() # 自动开始调度
| 方法 | 参数 | 说明 |
|---|---|---|
srv.producer(topic, interval, serializer, compression) |
interval 秒 |
固定间隔定时推送 |
srv.burst_producer(topic, serializer, compression) |
无间隔 | 连续推送,返回 None 停止 |
回调返回值支持 DataFrame / dict / str / bytes,自动编码为协议帧并路由到所有匹配的消费者。
生产者
import asyncio
from pulsemq.client import ProducerClient
async def main():
prod = ProducerClient(
"tcp://127.0.0.1:5555", "tcp://127.0.0.1:5556",
username="user1", password="pass1",
)
await prod.start()
await prod.publish("market.stock.AAPL", {"price": 180.5, "volume": 1000})
await prod.stop()
asyncio.run(main())
消费者
from pulsemq.client import ConsumerClient
async def main():
cons = ConsumerClient(
"tcp://127.0.0.1:5555", "tcp://127.0.0.1:5556",
username="user2", password="pass2",
)
await cons.start()
def on_msg(msg):
print(msg.topic, msg.payload, msg.timestamp_ns)
await cons.subscribe("market.*", on_msg)
await asyncio.sleep(3600) # 持续接收
asyncio.run(main())
架构
┌──────────────────┐
│ Server │
│ ┌─ 数据面 ─────┐│──── DEALER → 消费者
│ │ ROUTER :5555 ││
│ └──────────────┘│
│ ┌─ 控制面 ─────┐│
│ │ ROUTER :5556 ││←── DEALER ← 生产者
│ └──────────────┘│
│ ┌─ Admin ──────┐│
│ │ HTTP :9090 ││─── REST / SSE / Web UI
│ └──────────────┘│
│ ZAP (bcrypt) │
└──────────────────┘
| 端口 | 协议 | 用途 |
|---|---|---|
5555 |
ROUTER (数据面) | 消息发布/接收 |
5556 |
ROUTER (控制面) | REGISTER / HEARTBEAT / SUBSCRIBE / DISCONNECT |
9090 |
HTTP | 监控 Web UI + REST API + SSE |
数据流
生产者 DEALER ──encode→ ROUTER decode_header match topic ──→ 消费者 DEALER
↓
TrafficStats.record LatencyStats.sample
服务端不解压/不反序列化 payload(decode_header 仅提取头部),转发后由消费者完整 decode 还原。
控制面
| 命令 | 方向 | 作用 |
|---|---|---|
REGISTER |
Client → Server | 注册上线(含用户名、角色、订阅列表) |
HEARTBEAT |
Client → Server | 保活(每秒 1 次,6 秒超时自动下线) |
SUBSCRIBE |
Client → Server | 订阅 topic 模式 |
UNSUBSCRIBE |
Client → Server | 取消订阅 |
DISCONNECT |
Client → Server | 优雅下线 |
数据类型与序列化
支持的数据类型
| Python 类型 | 可用序列化器 | record_count |
|---|---|---|
pd.DataFrame |
msgpack, json, pyarrow | 行数 |
dict |
msgpack, json, pyarrow | 1 |
str |
str(仅此一种) | 1 |
bytes |
bytes(仅此一种) | 1 |
pyarrow 对 DataFrame/dict 均可直接序列化;json/msgpack 下 DataFrame 先转
list[dict]再以data_type=DATAFRAME标记,接收端_restore_type还原回 DataFrame。
序列化格式
| 格式 | 后端 | 适合 | 批处理场景 |
|---|---|---|---|
msgpack |
msgspec | 结构化小消息 ✅ | 批量 DataFrame 需先 to_dict |
json |
msgspec | 人类可读、跨语言 | 同上 |
pyarrow |
pyarrow IPC | 列存/分析 ✅ | 直接序列化 DataFrame,最快 |
str |
UTF-8 | 纯文本 | ❌ |
bytes |
透传 | 二进制 | ❌ |
压缩算法
| 算法 | 适用场景 |
|---|---|
none |
小消息,极速 |
snappy |
速度优先 |
lz4 |
批数据,平衡 |
zstd |
压缩比优先(带宽受限) |
小消息场景压缩是负收益(计算开销 > 传输节省);批量 DataFrame 场景 lz4/zstd 有明显效果。
客户端生命周期
首次启动
| 场景 | 异常 | exit code |
|---|---|---|
| 密码错误 | AuthenticationError |
3 |
| 服务器不可达 | ClientStartupError |
4 |
运行期重连
断线后 Client._reconnect_loop 按指数退避自动重连:
断线 → disconnected → cancel bg tasks → 新 Transport → PLAIN 认证
→ REGISTER(同 client_id)
├─ ALREADY_ONLINE → 退避重试(等待心跳超时释放)
├─ auth_failed → _reconnect_fatal → exit 3
└─ OK → 恢复订阅 → 重启 recv/heartbeat
业务无感:订阅自动恢复,消息继续接收。
监控
Web UI
浏览器打开 http://localhost:9090/?token=<admin_token> 查看实时面板:
- 4 个指标卡片:活跃主题、消息量/秒、流量/秒、运行时间
- 4 个客户端卡片:在线用户、生产者数、消费者数、订阅数
- ECharts 流量趋势折线图(分钟级,1H/6H 切换,最多 5 topic 叠加)
- 延迟 P50/P95/P99 柱状图
- 实时事件流(认证/连接/断线)
- 在线 Client 详情弹窗
REST API
# 实时指标
curl 'http://localhost:9090/api/v1/stats/realtime?token=<token>'
# 主题列表
curl 'http://localhost:9090/api/v1/topics?token=<token>'
# 主题分钟级历史
curl 'http://localhost:9090/api/v1/topics/market.tick/history?minutes=60&token=<token>'
# 在线客户端明细
curl 'http://localhost:9090/api/v1/clients?token=<token>'
# 生命周期事件
curl 'http://localhost:9090/api/v1/events?token=<token>'
# 健康检查(无需 token)
curl http://localhost:9090/healthz
SSE 实时流
curl -N 'http://localhost:9090/api/v1/stats/stream?token=<token>'
每 1 秒推送一帧 JSON,包含 topics / latency / online_users / sse_events 等。
协议帧格式
单 bytes 帧(非 ZMQ 多帧,通过 DEALER/ROUTER 传输):
magic(2) ver(1) msg_type(1) flags(1) data_type(1) topic_len(2 BE)
topic(N) ts(8 BE ns) record_count(4 BE) payload(变长) [CRC32?(4)]
magic="PM"msg_type= DATA(0x01) / CONTROL(0x02)flags= 编码序列化器(3bit) + 压缩算法(2bit) + CRC(1bit)data_type= UNKNOWN(0x00) / DICT(0x01) / DATAFRAME(0x02) / STR(0x03) / BYTES(0x04)record_count上限 1,000,000- CRC 可选(由 flags 指示)
凭据管理
CLI
# 添加用户(自动 bcrypt 哈希)
pulsemq users add trader1 secret123 --role publisher --role subscriber
# 启动服务端
pulsemq server
# 列出所有用户
pulsemq users list
# 禁用/启用用户
pulsemq users disable trader1
pulsemq users enable trader1
# 修改密码
pulsemq users passwd trader1 new_secret
# 热加载凭据(SIGHUP 或 CLI)
pulsemq users reload
文件格式 (pulsemq_users.toml)
[users.admin]
hashed_password = "$2b$12$..."
roles = ["admin"]
enabled = true
created_at = "2026-06-27T00:00:00Z"
Admin Token
首次启动自动生成 32 字节随机 base64url token,写入 pulsemq_admin.token(0600 权限)。
Web UI 和 REST API 通过 ?token=... 或 Authorization: Bearer ... 传递。
可通过环境变量 PULSEMQ_ADMIN_TOKEN 或配置文件覆盖。
日志
日志输出到 logs/ 目录,每日滚动,保留 30 天:
logs/
├── pulsemq_2026-06-27.log
├── pulsemq_2026-06-28.log
└── ...
stderr 同步输出(容器/交互可见)。
性能基准
单条消息(dict,每帧 1 条)
| 序列化 | 压缩 | 吞吐量 | P50 延迟 |
|---|---|---|---|
| msgpack | none | 18.9 K/s | 183.9 ms |
| msgpack | lz4 | 18.4 K/s | 124.1 ms |
| json | none | 18.6 K/s | 151.4 ms |
| pyarrow | none | 2.1 K/s | 781.6 ms |
pyarrow 适用于批处理,单条 dict 场景序列化开销过大(8-9x 慢于 msgpack)。 压缩对< 200B payload 为负收益。机器:Windows 11, Python 3.14, 单机 localhost。
批量 DataFrame(1000 行/帧)
| 序列化 | 压缩 | 记录/s | P50 延迟 |
|---|---|---|---|
| pyarrow | zstd | 363.6 K/s | 1677 ms |
| pyarrow | lz4 | 341.1 K/s | 1481 ms |
| msgpack | none | 292.5 K/s | 1955 ms |
| json | none | 290.5 K/s | 1983 ms |
批量场景 pyarrow + 压缩是最优组合(直接序列化 DataFrame,无需格式转换)。 延迟为压测下 ZMQ 缓冲区排队所致,真实场景以固定间隔发送时远低于此。
基准脚本
# 端到端压测(单进程 consumer + producer)
python bench_v2_e2e.py --duration 10
# 分离进程行情压测
python bench_v2_market.py producer --duration 10
python bench_v2_market.py consumer --duration 13
# 全矩阵 ser × comp × data_type 类型保真验证
python bench_v2_matrix.py
# 行情全矩阵性能(12 组合)
python bench_market_full.py --duration 5
# DataFrame 批量性能(1000 行/帧)
python bench_df_batch.py --duration 5
配置
环境变量
| 变量 | 说明 | 默认 |
|---|---|---|
PULSEMQ_DATA_ENDPOINT |
数据面绑定地址 | tcp://0.0.0.0:5555 |
PULSEMQ_CONTROL_ENDPOINT |
控制面绑定地址 | tcp://0.0.0.0:5556 |
PULSEMQ_ADMIN_ENDPOINT |
管理 HTTP 绑定地址 | 0.0.0.0:9090 |
PULSEMQ_CREDENTIALS_FILE |
凭据 TOML 路径 | ./pulsemq_users.toml |
PULSEMQ_ADMIN_TOKEN |
监控接口 token | 自动生成 |
PULSEMQ_ADMIN_TOKEN_FILE |
token 文件路径 | ./pulsemq_admin.token |
PULSEMQ_ADMIN_PASSWORD |
首次启动默认密码 | 随机 16 字符 |
PULSEMQ_STATS_DB |
SQLite 路径 | ./pulsemq_stats.sqlite |
PULSEMQ_HEARTBEAT_TIMEOUT |
心跳超时(秒) | 6.0 |
PULSEMQ_LATENCY_SAMPLE_RATE |
延迟采样率 | 0.01 (1%) |
PULSEMQ_STATS_RETENTION |
内存窗口(分钟) | 480 (8h) |
PULSEMQ_BCRYPT_COST |
bcrypt 代价因子 | 12 |
更新日志
v2 (current)
完整重构:PUB/SUB → Client/Server (ROUTER/DEALER)
- 架构变更:单 PUB socket → 双 ROUTER(数据面 + 控制面)+ HTTP admin
- 认证升级:api_key 明文 → bcrypt CredentialStore + ZAP PLAIN + 用户 CLI
- 类型保真:
_restore_type确保 DataFrame/dict str/bytes 端到端还原 - 监控增强:延迟 P50/P95/P99、在线客户端、事件流、SSE 实时推送、独立 Admin 线程
- 性能优化:
decode_header服务端零反序列化路由、TrafficStats单dict.get - 自动重连:运行期指数退避重连(1s→2s→4s→...→30s)
- 日志系统:loguru 统一,每日滚动写入
logs/,30 天保留 - 安全性:密码 bcrypt 哈希、admin token 随机生成(0600)、凭据文件原子写入
- 批次处理:DataFrame 批量 1000 行/帧 达到 363 K records/s(pyarrow+zstd)
许可证
Project details
Release history Release notifications | RSS feed
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file pulse_mq-6.0.0.tar.gz.
File metadata
- Download URL: pulse_mq-6.0.0.tar.gz
- Upload date:
- Size: 596.7 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.14.6
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
b7c39147c76c14974c214090599cb427f57a43dcdb12717d80f8bb5228bedc7d
|
|
| MD5 |
564d0c2d181bb9d3be3a7378bbe44f31
|
|
| BLAKE2b-256 |
40635ad933e25756849763db3de1c8c3454dc5b17e6ff120c586bc77ccbad39d
|
File details
Details for the file pulse_mq-6.0.0-py3-none-any.whl.
File metadata
- Download URL: pulse_mq-6.0.0-py3-none-any.whl
- Upload date:
- Size: 418.7 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.14.6
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
340b4fc325f4ee8b7629f8cb96582f256b67283f5eb4a053a6e563fe7331c613
|
|
| MD5 |
38a600bf41927f16cbcebc8b1e7a5cc3
|
|
| BLAKE2b-256 |
bdd20e3a07a49751d2d204469200a808c5ae4d0af488c46d0d3524708124280d
|