这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 会用 API 只是开始,把 Kafka 放进架构需要处理一致性、顺序、依赖与演进。本章讲事件驱动架构中的核心模式。
26.1 事件驱动 vs 请求驱动
| 维度 | 请求驱动(RPC/REST) | 事件驱动(Kafka) |
|---|---|---|
| 耦合 | 调用方必须知道下游 | 生产者只发布事实 |
| 时间 | 下游必须在线 | 异步,可落后再追 |
| 扩展 | 新下游要改调用方 | 新订阅者零侵入 |
| 一致性 | 强(同步返回) | 最终一致 |
| 排查 | 链路追踪直观 | 需要 event id + 对账 |
判断标准:“这个动作是否是已经发生的事实?” 是则发布事件;“希望对方做什么”更适合命令/请求。
26.2 三种事件风格
1. Event Notification(事件通知)
OrderCreated { orderId }
只告诉“发生了”,下游需要回查。优点:事件小、敏感数据少;缺点:下游强依赖上游查询接口,形成隐式耦合。
2. Event-Carried State Transfer(携带状态)
OrderCreated { orderId, amount, userId, status }
事件自带数据,下游可本地物化,真正解耦。代价:数据冗余与 schema 演进治理(第 8 章)。
3. Event Sourcing(事件溯源)
事件是唯一事实源,当前状态由事件回放得到:
事件流: Created -> ItemAdded -> Paid -> Shipped
当前状态: 由 1..n 折叠(fold)得到
适合审计、回放、复杂状态机;成本高(schema 稳定性、快照、版本迁移),不要无脑上。
26.3 Outbox 模式:终结双写
问题:业务写库和发 Kafka 是两个系统,无法放进一个事务。
错误做法:
BEGIN; 写订单; COMMIT; sendKafka(...); // 中间宕机 -> 丢事件
sendKafka(...); BEGIN; 写订单; COMMIT; // 发了事件但库没写 -> 幽灵事件
Outbox 解法:
本地事务:
INSERT orders;
INSERT outbox(event);
COMMIT;
异步投递:
relay 读 outbox -> 发 Kafka(至少一次)
消费端按 event_id 幂等
实现细节:
- outbox 表加
status/created_at/id,Relay 按 id 顺序投递并标记; - 投递用
acks=all+ 回调确认; - 允许重复发送,幂等在消费端兜底;
- 表数据量大时按时间分区/定期归档。
26.4 Saga:跨服务最终一致
没有分布式事务时,用“本地事务 + 补偿”编排长流程:
下单 Saga:
1. 订单服务: 创建订单(PENDING) [本地事务]
2. 库存服务: 预扣库存 [本地事务 + 消费事件]
3. 支付服务: 创建支付单
4. 成功: 订单 PAID, 库存 CONFIRMED
失败: 订单 CANCELLED, 库存 RELEASE(补偿)
两种编排方式:
| 方式 | 特点 |
|---|---|
| 编排(Orchestration) | 中央 Saga 协调器发命令、监听结果,流程清晰 |
| 协同(Choreography) | 服务互相订阅事件,去中心化,链路难追踪 |
实践建议:超过 3-4 步、有超时回滚的流程用编排;简单链路用协同 + 全链路 event id。
26.5 CQRS 与 Kafka
命令写主库,查询读物化视图:
Command -> MySQL -> outbox -> Kafka ->
Projector(消费者) -> ES/Redis/ClickHouse 读模型
要点:
- 读模型按查询场景建(宽表、索引、物化列);
- 投影是幂等的,重启可从 offset 重放;
- 最终一致窗口要监控(主库 vs 读模型延迟);
- 业务要能接受“写后读可能读旧值”(或读主库兜底)。
26.6 背压
Kafka 天然是缓冲区,但不是无限缓冲:
洪峰 -> Kafka 堆积(lag) -> 消费者按能力拉取 -> 下游稳定
背压治理三问:
- 堆积多久可接受?(业务 SLA)
- 保留期能否覆盖最长堆积?(防位移过期)
- 下游能否弹性扩容?(消费者数 <= 分区数是硬上限)
超过分区数还想加速:分区扩容(注意 key 顺序)、下游再并行、或拆分 topic。
26.7 多租户隔离
常见三级隔离:
| 级别 | 做法 | 隔离度 |
|---|---|---|
| 共享 topic | 租户 ID 做消息字段/header | 弱,易互相影响 |
| 独立 topic | 每租户一组 topic + ACL | 中,配额与监控友好 |
| 独立集群 | 物理隔离 | 强,成本高 |
配套手段:
- quota 按用户/客户端限流,防止单租户打挂集群;
- ACL 严格隔离读写;
- header 放租户 ID,消息体可加密;
- 大租户独立 topic/分区,避免热点 key 倾斜。
26.8 Schema 与契约演进
事件是跨团队契约,演进要有流程:
1. schema 入 Git,PR 评审
2. CI 跑兼容性检查(FULL 推荐)
3. 消费者先兼容旧版,再发布新版
4. 破坏性变更 -> 新 topic / 新版本事件
5. 双写过渡期 -> 验证 -> 下线旧版
配套治理:
- event id 全局唯一,贯穿链路追踪与对账;
- 事件带
occurred_at(业务时间)与version; - 建立 schema registry 与消费关系登记,知道“谁在用这个事件”。
26.9 事件驱动架构的坑
- 把事件当命令用:到处发布
DoSomething,职责倒挂; - 无幂等:重复消费直接写脏数据;
- 无对账:最终一致变成“最终不知道一不一致”;
- 过度事件溯源:普通业务上了全套,复杂度爆炸;
- 共享主题热 key:大客户流量把一个分区打满;
- 无死信治理:失败消息进黑洞;
- 版本随意变更:下游大规模反序列化失败。
本章小结
- 事件表达“已发生的事实”,命令表达“要求执行的动作”,两者别混;
- Outbox 解决双写,Saga 解决跨服务长事务,CQRS 解决读写模型差异;
- 幂等、对账、死信治理是事件系统的三大基础设施;
- 背压是 Kafka 的天赋,但保留期与分区数是硬边界;
- schema 是团队间契约,演进必须流程化。
思考题
- “PaymentRequested” 和 “PaymentCompleted” 哪个是事件?哪个更像命令?为什么?
- Saga 中补偿动作本身失败怎么办?设计重试与人工介入机制。
- 你的系统里哪些查询适合走 CQRS 读模型?数据延迟窗口要求是多少?