14、消息队列与事件驱动:Kafka/RabbitMQ选型、事件溯源、CQRS模式、异步处理
做市商系统里,最怕什么?
我个人最怕的就是「耦合」。报价模块挂了,把交易模块拖死;成交回报慢了,把风控模块堵死。说白了,一个环节出问题,整条链路跟着遭殃。
怎么解?消息队列 + 事件驱动。这套组合拳,是我在搭建结构化产品做市系统时,花了大半年才打磨顺手的。今天聊聊我的实战心得。
14.1 为什么做市商系统离不开消息队列?
先想一个问题:行情来了,每秒几千笔报价更新,你的系统怎么处理?
同步调用?那完了。行情模块等着风控模块算完,风控等着定价引擎算完,定价引擎又等着数据源更新……这一圈下来,行情早过期了。
异步才是正解。消息队列就是那个「中间人」——行情来了,丢进队列,谁有空谁处理。处理慢了?队列帮你缓冲。处理挂了?消息还在,重启后继续消费。
我在项目中遇到过最典型的情况:某次行情暴涨,报价量翻了10倍。如果没有消息队列做削峰填谷,数据库直接被打满,整个系统瘫痪了20分钟。从那以后,我对消息队列的依赖就再也没动摇过。
14.2 Kafka vs RabbitMQ:选型对比
选哪个?这问题我被人问了不下50次。直接上对比表:
| 维度 | Kafka | RabbitMQ |
|---|---|---|
| 设计定位 | 高吞吐、持久化、流处理 | 可靠路由、灵活交换、低延迟 |
| 吞吐量 | 百万级消息/秒 | 万级消息/秒 |
| 消息模型 | 拉模式(Consumer Pull) | 推模式(Push) + 拉模式 |
| 消息顺序 | 分区内严格有序 | 单队列有序,多队列无序 |
| 消息回溯 | 支持(基于Offset) | 不支持(消费后删除) |
| 典型场景 | 事件溯源、日志收集、流计算 | 任务分发、RPC解耦、事务消息 |
我的建议很简单:
- 做市商核心链路(行情、订单、成交) → 选Kafka。吞吐量大,消息能回溯,适合事件溯源。
- 管理后台、通知、任务调度 → 选RabbitMQ。路由灵活,延迟低,运维简单。
嗯,这里要注意:别想着一个队列打天下。我见过有人用RabbitMQ扛行情,结果消息堆积到几百万,直接OOM。也见过用Kafka做任务调度,延迟高得离谱。工具选对,事半功倍。
14.3 事件溯源:把「状态」变成「事件流」
传统做法是存「当前状态」——订单表里一条记录,状态是「已成交」。但问题来了:如果我想知道这个订单是怎么一步步变成「已成交」的?中间经历了哪些状态变更?谁在什么时间修改了它?
事件溯源的做法是:不存状态,只存事件。
举个例子:
// 传统方式
Order {
id: "ORD001",
status: "FILLED",
filledQty: 100,
updateTime: "2024-01-15 10:30:00"
}
// 事件溯源方式
Event 1: OrderCreated { id: "ORD001", qty: 100, price: 105.20, time: "10:29:00" }
Event 2: OrderPartiallyFilled { id: "ORD001", filledQty: 50, time: "10:29:30" }
Event 3: OrderFilled { id: "ORD001", filledQty: 50, time: "10:30:00" }
你想想看,事件溯源的好处是什么?
- 完整的审计轨迹:每个状态变更都有记录,谁改的、什么时候改的、改之前是什么样,一清二楚。
- 时间旅行:我可以把系统回滚到任意时间点,重现当时的场景。做市商系统做回测时,这个能力太重要了。
- 天然支持CQRS:事件是写模型,读模型可以从事件流中任意构建。
避坑指南:我曾经在事件溯源中犯过一个低级错误——事件结构定义得太死板。后来业务需求变了,事件字段不够用,只能做事件版本升级,那叫一个痛苦。建议从一开始就给每个事件加上 version 字段,预留扩展空间。
14.4 CQRS模式:读写分离的终极形态
CQRS(Command Query Responsibility Segregation),说白了就是「写操作走一条路,读操作走另一条路」。
为什么做市商系统需要CQRS?
- 写操作:下单、撤单、改单,要求强一致性、低延迟。通常走Kafka,保证顺序和持久化。
- 读操作:查询持仓、查询历史成交、生成报表,要求高并发、灵活查询。通常走Redis或Elasticsearch,甚至可以是物化视图。
我画了一张图,帮你理解CQRS在结构化产品做市系统中的位置:
这张图里,核心逻辑是:
- 命令端只管「写」,把事件丢进Kafka
- 事件存储把所有事件持久化,形成完整的事件流
- 投影(Projection)从事件流中构建出读模型,存到Redis或ES里
- 查询端只管「读」,从读模型里拿数据,性能极高
小技巧:读模型可以不止一个。比如持仓查询用一个读模型,风控指标用另一个读模型,报表统计再用一个。每个读模型只关注自己需要的事件,互不干扰。我在项目中最多维护了7个不同的投影,各自独立更新,效果很好。
14.5 异步处理:别让「等待」成为瓶颈
做市商系统里,异步处理无处不在。我总结了几种典型模式:
14.5.1 异步下单
用户下单后,系统立刻返回「已接收」,然后异步去校验风控、检查资金、撮合交易。用户不用傻等。如果后面出问题了,通过回调或轮询通知用户。
14.5.2 异步风控
每笔订单进来,先丢进Kafka,风控模块异步消费。风控规则复杂?没关系,慢慢算。算完了发个事件,订单模块根据结果决定是否继续。
14.5.3 异步报表
日报、周报、月报,这些计算量大的任务,全部异步处理。我习惯用RabbitMQ做任务分发,多个worker并行计算,最后汇总。
警告:异步不是银弹。异步带来的问题包括:
- 最终一致性:写和读之间有时延,用户可能看到「过时」的数据。需要业务上容忍。
- 消息丢失:Kafka虽然持久化,但配置不当还是会丢消息。记得开启 acks=all 和 min.insync.replicas=2。
- 重复消费:网络抖动可能导致消息重复。消费端要做好幂等处理。
我曾经因为没做幂等,导致一笔订单被重复成交了两次,亏了十几万。从那以后,每个消费者我都强制加上去重逻辑。
14.6 实战:一个完整的异步事件流
最后,给你看一个我在结构化产品做市系统中实际用过的异步事件流:
// 1. 用户下单 -> 产生 OrderPlaced 事件
Event: OrderPlaced {
orderId: "ORD-20240115-001",
productId: "SNOWBALL-001",
side: "BUY",
quantity: 100,
price: 105.20,
timestamp: 1705310400000
}
// 2. Kafka 将事件持久化,并分发给多个消费者
// 消费者A:风控校验
// 消费者B:资金检查
// 消费者C:行情快照
// 3. 风控通过 -> 产生 RiskPassed 事件
Event: RiskPassed {
orderId: "ORD-20240115-001",
checkResult: "PASS",
timestamp: 1705310400100
}
// 4. 资金充足 -> 产生 FundReserved 事件
Event: FundReserved {
orderId: "ORD-20240115-001",
reservedAmount: 10520.00,
timestamp: 1705310400200
}
// 5. 两个事件都到达后 -> 订单进入撮合引擎
// 撮合引擎消费这两个事件,确认条件满足,开始撮合
// 6. 撮合完成 -> 产生 OrderExecuted 事件
Event: OrderExecuted {
orderId: "ORD-20240115-001",
executedQty: 100,
avgPrice: 105.18,
timestamp: 1705310400500
}
// 7. 投影更新读模型
// 持仓读模型:增加100股
// 资金读模型:扣减10518元
// 报表读模型:记录一笔成交
你看,整个过程全是事件驱动,没有一次同步调用。每个模块各司其职,互不阻塞。这就是消息队列 + 事件驱动的魅力。
嗯,最后说一句:这套架构不是一天建成的。我建议你先从最核心的链路开始,比如行情 -> 报价 -> 成交,用Kafka串起来。跑顺了,再逐步扩展到风控、资金、报表。步子迈大了,容易扯着蛋。
无相订单流研究社 微信Lucian808555