{T}

Kafka 与 ZMQ 消息队列

当量化交易系统从单机走向分布式时,模块间的通信方式就成为了架构的核心问题。消息队列是解决这个问题的标准答案——它实现了模块间的解耦异步可靠通信。

阅读提示

为什么需要消息队列?

从单机到分布式

图表渲染中…

单机模式下,模块之间直接调用。分布式模式下,模块分散在不同机器上,需要消息队列作为通信中枢。

中间件的概念

中间件是连接底层技术和上层应用的组件——让使用者不必关心底层实现,只需调用统一接口

图表渲染中…

中间件哲学:没有什么事情是加一层解决不了的;如果有,那就加两层。

消息队列的核心特性

发布-订阅模式(Pub-Sub)

图表渲染中…
特性说明
时序性严格 FIFO(先进先出),丢入 1, 2, 3,取出也是 1, 2, 3
解耦发布者和订阅者之间无直接依赖,互相不知道对方的存在
可靠性内置消息持久化、去重和重试机制
可扩展新增订阅者不影响现有系统和业务

ZMQ:轻量级消息队列

ZMQ(ZeroMQ)是一个高性能异步消息库,提供三种通信模式:

模式说明适用场景
REQ-REP请求-应答RPC 调用,类似 HTTP
PUB-SUB发布-订阅行情数据广播
PUSH-PULL推送-拉取并行任务分发

ZMQ PUB-SUB 实战

python
# ============ 订阅者 1 (sub1.py) ============
import zmq


def run():
    context = zmq.Context()
    socket = context.socket(zmq.SUB)
    socket.connect('tcp://127.0.0.1:6666')
    # 空字符串 = 不过滤任何消息,接收全部
    socket.setsockopt_string(zmq.SUBSCRIBE, '')

    print('订阅者 1 已启动')
    while True:
        msg = socket.recv()
        print(f"订阅者 1 收到: {msg}")


if __name__ == '__main__':
    run()


# ============ 订阅者 2 (sub2.py) ============
import zmq


def run():
    context = zmq.Context()
    socket = context.socket(zmq.SUB)
    socket.connect('tcp://127.0.0.1:6666')
    # 只订阅以 'BTC' 开头的消息
    socket.setsockopt_string(zmq.SUBSCRIBE, 'BTC')

    print('订阅者 2 已启动 (仅接收 BTC)')
    while True:
        msg = socket.recv()
        print(f"订阅者 2 收到: {msg}")


if __name__ == '__main__':
    run()


# ============ 发布者 (pub.py) ============
import time
import zmq


def run():
    context = zmq.Context()
    socket = context.socket(zmq.PUB)
    # 注意:PUB 用 bind,SUB 用 connect
    # 同一个地址端口 bind 只能有一个,connect 可以有多个
    socket.bind('tcp://*:6666')

    cnt = 1
    while True:
        time.sleep(1)
        # 发送 BTC 行情
        socket.send_string(f'BTC price:{50000 + cnt * 100}')
        # 发送 ETH 行情
        socket.send_string(f'ETH price:{3000 + cnt * 10}')
        print(f'发布者发送第 {cnt} 条')
        cnt += 1


if __name__ == '__main__':
    run()

运行顺序

code
1. 先启动所有订阅者 (sub1.py, sub2.py)
2. 再启动发布者 (pub.py)

为什么必须先启动订阅者? ZMQ 的 PUB-SUB 模式中,订阅者连接到发布者时需要时间来建立连接。如果先启动发布者,早期的消息会在订阅者连接之前发送出去,导致消息丢失。

ZMQ 关键设计点

设计决策原因
PUB 用 bind,SUB 用 connectbind 只能有一个(发布者独占端口),connect 可以有多个
setsockopt_string(zmq.SUBSCRIBE, '')空字符串 = 接收所有消息;指定前缀 = 按主题过滤
socket.recv() 是阻塞的没有新消息时,订阅者会阻塞等待

Kafka:工业级消息队列

Kafka 核心概念

图表渲染中…
概念说明
Topic消息的逻辑分类(如 market-dataorders
Producer消息生产者,向 Topic 发送消息
Consumer消息消费者,从 Topic 读取消息
Consumer Group消费者组,组内消费者分摊 Partition
BrokerKafka 服务器节点
PartitionTopic 的分片,实现并行处理和水平扩展

ZMQ vs Kafka 对比

维度ZMQKafka
本质消息库(library),嵌入应用进程消息平台(platform),独立部署的集群
持久化无内置持久化(消息在内存中)磁盘持久化,可回溯历史消息
吞吐量极高(百万条/秒)
运维复杂度低(只需安装 Python 包)高(需部署和管理 Kafka 集群)
适用场景实时数据流、低延迟广播大规模数据管道、事件溯源、日志收集
消息回溯不支持支持(可按 offset 重新消费)
典型用户实时行情推送订单日志、用户行为分析

量化交易系统中的选择

图表渲染中…
场景选择原因
实时行情推送ZMQ低延迟、轻量级、部署简单
订单指令传递Kafka必须持久化、不可丢失、可回溯
日志收集Kafka大量数据、需要持久化和后续分析
模块间 RPC 调用ZMQ REQ-REP简单的请求-应答模式

分布式架构中的消息队列

Single Point Failure 问题

code
❌ 单机监听交易所的风险:

行情模块 (单节点) ──── 交易所
        ↓
   这台机器故障/网络中断
        ↓
   整个交易系统暴露于风险中

多机冗余方案

图表渲染中…

优点:消除 Single Point Failure 缺点:增加数十到数百毫秒延迟(等待多个节点数据 + 消息队列处理) 适用:低频策略、波段策略(延迟换稳定性);不适用高频策略

常见问题

Q: 把 time.sleep(1) 放在 while 循环最后 vs 最前面有区别吗?

A: 放在循环最后,第一次迭代时没有延迟,会立即发送第一条消息。放在循环最前面,第一次迭代会先 sleep 1 秒再发送。对于订阅者来说,放在最后意味着第一条消息更快到达。

Q: 多个发布者怎么用 ZMQ?

A: 每个发布者 bind 到不同的端口,订阅者连接到多个端口。或者使用 ZMQ 的 proxy 模式(XPUB/XSUB)来汇聚多个发布者的消息。

Q: ZMQ 和 RabbitMQ 有什么区别?

A: RabbitMQ 是一个完整的消息代理(broker),需要独立部署和运维。ZMQ 只是一个库,直接嵌入应用进程,没有中心化的 broker。ZMQ 更轻量但功能更少,RabbitMQ 更重但提供更多企业级特性(管理界面、消息确认、持久化等)。

术语表

术语英文定义
消息队列Message Queue临时存放消息的容器,支持发布-订阅、点对点等通信模式
Pub-SubPublish-Subscribe发布-订阅模式:发布者发送消息到主题,订阅者接收感兴趣的消息
FIFOFirst In First Out先进先出,队列的核心顺序保证
Single Point FailureSPOF单一故障点:系统中的一个组件故障会导致整个系统不可用
中间件Middleware连接底层技术和上层应用的组件层

延伸阅读