Chapter 05

综合实战:订单事件系统的选型与判别

01 章把 Kafka 框成一条可重放的分区日志,02 章拆开了副本/再平衡/投递语义的设计取舍,03 章把它写成可跑的代码,04 章列出机制被违反时的失败模式。本章不再逐节讲机制,而是把前四章塞进一个真实系统:一个电商订单事件平台。它逼出的不是"怎么实现某步",而是判别——同一个决策点,这里该用 01 的概念还是 02 的机制,该选 Kafka 还是别的工具。

本章你将建立的 schema

  • 把"一份数据、多消费组、各自 offset"从概念落到一张订单系统架构图上
  • 分区 key、acks 组合、消费语义、顺序范围四类决策——每类都是一次跨章判别,没有放之四海的默认答案
  • 选型判别:同一个系统里,为什么主干用 Kafka、某个下游反而该用 RabbitMQ、什么时候考虑 share group
  • 用可观测行为(重启不重复扣、kill broker 不丢、同单严格有序)来验收设计,而不是"看起来对"
订单服务 key = orderId TOPIC orders(3 分区) partition 0 · 日志 partition 1 · 日志 partition 2 · 日志 写 消费组 · 库存扣减 offset 库存 · 独立游标 不丢且不重复扣 消费组 · 用户通知 offset 通知 · 独立游标 可容忍少量重复 消费组 · 财务对账 offset 对账 · 独立游标 不丢 + 按订单有序 各读全量 读不删数据 → 可重放
图 5.1一份订单日志,三个消费组各读各的、各存各的 offset。 注意:三条出边都从同一个 orders topic 出发——数据不被复制三份,三组各自持有独立游标(这正是 01 章"offset 在消费者侧")。三组对"丢失/顺序/重复"的要求不同(朱红标注),所以下面每个决策都要分组回答,而不是一刀切。

5.1项目背景:订单事件平台

一个电商平台的订单服务,每当一笔订单发生状态变化就产生一个事件:订单创建、订单支付、订单发货。这些事件统一写进一个 orders topic。三个互不相关的团队各自起一个消费组来读这条流:

  • 库存扣减(inventory):订单支付后扣减对应 SKU 库存。要求不丢(少扣会超卖)且不重复扣(重复扣会少卖、库存账错)。扣减动作落在一个外部库存数据库上。
  • 用户通知(notification):给用户发短信/推送"订单已支付""已发货"。要求不丢(用户该收到),但可容忍少量重复——偶尔多收一条推送是体验问题,不是数据正确性问题。
  • 财务对账(reconciliation):把订单事件写进对账系统、生成财务流水。要求严格不丢(少一笔对不平)且按订单顺序处理——同一笔订单的"创建→支付→发货"必须按发生顺序入账,否则状态机错乱。
量级与需求画像

吞吐:大促峰值约每秒数万条订单事件,平时每秒数千条——属于高吞吐事件流,远超"低吞吐小消息"的反场景。

顺序:只在单笔订单内部需要先后顺序(创建先于支付先于发货)。两笔不同订单之间没有顺序约束。没有任何下游需要"全平台所有订单的全局总序"。

可靠性:库存与对账是不能丢的强一致诉求;通知是尽量别丢、但重复无害的弱诉求。一份数据上,三组的可靠性等级不同。

重放:对账系统改了计算口径要重算上月流水、库存对错了要从某天重新推演——都需要把消费组 offset 重置后重放历史。这是选 Kafka(而非删除式队列)的硬需求之一。

这套画像把"选 Kafka"几乎写死了:高吞吐 + 一份数据多组独立消费 + 需要重放,正是 04 章表 4.1 的"选中"行。但"选了 Kafka"只是起点——真正的工程判断在下面四个决策点上,每一个都没有"标准答案",只有"对这一组、在这个量级下"的答案。

5.2设计任务:五个判别决策

下面每一行都是一次判别:给出备选、给出这个场景下的选择、给出"为什么不选另一条"。每个决策点回链到前面定义过它的章节——判别题考的从来不是记住配置名,是对取舍的推理。

表 5.1 · 订单系统的五个跨章判别决策
决策点备选这个场景选哪个 / 为什么
① 分区 key
01 §1.3 · 02 §2.2 · 04 §4.3
userId 做 key / orderId 做 key / 不带 key 选 orderId。同单事件靠相同 key 落同一分区从而有序(满足对账的按订单有序);orderId 高基数,分布均匀。userId 会让大客户的所有订单挤进一个分区造成热分区倾斜;不带 key 则同单事件会按 sticky partitioner 散到不同分区、丢掉顺序保证。
② 持久性组合
02 §2.3 · 04 §4.1
acks=1 / acks=all 单独 / acks=all+min.insync=2+RF=3+unclean=false 对账与库存选四件套组合。只设 acks=all 会在 ISR 缩到 1 时静默退化成 acks=1;持久性的真实来源是 min.insync.replicas。设 min.insync=RF-1=2:挂一台仍可写、挂两台拒写。unclean.leader.election=false 堵住落后副本上位截断已提交记录。
③ 库存的消费语义
02 §2.5 · 04 §4.2
at-least-once + 消费侧幂等 / 事务 EOS(read_committed) 选 at-least-once + 消费侧幂等。库存扣减的副作用落在外部数据库,而 Kafka 事务的精确一次只在 Kafka 闭环内成立、覆盖不到 DB 写。用幂等键(orderId+事件类型)做 upsert/唯一约束去重,比上事务更简单、且真正挡住重复。事务 EOS 在这里既增延迟又解决不了外部副作用。
④ 对账的顺序范围
01 §1.3 · 02 §2.2
全局严格顺序(单分区)/ 按订单顺序(分区内有序) 选按订单顺序。业务只要求"同一笔订单的事件有序",不需要跨订单全局序。靠 ② 的 orderId key 已经做到分区内有序。强上全局顺序要把 topic 压成单分区,并行度归 1、吞吐被杀死——为一个业务不需要的保证付全部吞吐的代价。
⑤ 选型 / 消费模型
04 §4.4
整体 Kafka / 某下游换 RabbitMQ / 普通 consumer group vs share group 主干选 Kafka(高吞吐+多组+重放)。但若通知组后续要按通知渠道做复杂逐条路由、要原生重试/死信队列,那一段反而更适合 RabbitMQ。若通知组要消费者数超过分区数、且每条独立 ack,可考虑 share group(KIP-932,4.2 GA)。详见 §5.3。

表里五行有一个共同点值得停下来看:同一个系统、同一份数据,三个消费组的最优答案不一样。对账要四件套 + 严格按序,通知可以松到能丢一点点配置、甚至换个 broker。把"全系统一套配置"当默认,是这道题最常见的错——判别的核心就是认出哪一组该用哪条机制。

想一想

决策 ① 选了 orderId 做 key。如果某天发现 orders 分区不够、想从 3 个加到 6 个,对"同一订单事件有序"这个保证会发生什么?

展开答案(先停 10 秒再点)

会被破坏。分区路由是 murmur2(orderId) % 分区数,把分区数从 3 改成 6,% 3 变 % 6,已有订单的 key→分区映射全部改变。一笔正在进行中的订单,它"创建/支付"时算出分区 1、加分区后"发货"算出分区 4——同一订单的事件分散到两个分区,分区内有序的前提没了,对账会出现先读到发货、后读到支付的乱序。

这正是 04 章"加分区静默破坏 key 顺序、且分区不可减"。正确做法:保守预估分区数,真要扩容就新建一个目标分区数的 topic 迁移,而不是原地加分区。这道题指向"分区数是早期就要算准的不可逆决策"。

5.3选型判别:Kafka / share group / RabbitMQ

决策 ⑤ 值得单独展开,因为它最考"是不是真懂 Kafka 的边界"。一个常见的红旗是把 Kafka 当"更快的队列",于是什么场景都往上套。04 章那张"什么时候不要用 Kafka"的表反过来定义了它的主场。下面这棵决策树把判别路径画出来——从负载特征一路走到具体选项。

选消息系统 需要可重放 + 高吞吐? 否 RabbitMQ 低量 / 路由 是 需要逐条路由 / 优先级 / 原生 DLQ? 是 RabbitMQ 该下游单拆 否 每条独立 ack / 消费者数 > 分区数? 是 share group KIP-932 · 4.2 否 consumer group 本系统主干
图 5.2消息系统选型决策树:三道菱形把负载特征筛到四个终点。 注意:终点不是"二选一",而是可以在一个系统里共存——订单主干走右下的 consumer group(朱红),通知这种带复杂路由的下游可以单独切到 RabbitMQ。"全用 Kafka"或"全不用"都是把判别题做成了站队题。

三道判别问各自在问什么

  • Q1 可重放 + 高吞吐? 这是 Kafka 的存在理由。订单系统要重放历史对账、峰值数万条/秒——是。低吞吐、不需要重放的小系统,分区/副本/再平衡的运维与认知成本换不回收益,传统队列更省。
  • Q2 逐条路由 / 优先级 / 原生 DLQ? Kafka 是分区日志,没有 exchange 式按内容路由、没有消息优先级、没有原生 per-message 重试与死信队列。通知组若要"按渠道路由 + 失败自动进 DLQ 重试",在 Kafka 上要自建重试 topic + 计数,RabbitMQ 原生就有——这一个下游单独用 RabbitMQ 是合理的混合架构,不是失败。
  • Q3 每条独立 ack / 消费者数 > 分区数? 普通 consumer group 的并行上限死等于分区数、且按分区批量推进 offset。share group(KIP-932,4.2 GA)提供队列语义:按记录 ack、消费者数不再受分区数限制。通知这种无序、想用很多消费者摊平瞬时积压的场景适合它;但它仍架在日志之上,不改变 Q1/Q2 的边界。
洞察 · 判别不是站队

资深信号不是"全用 Kafka"也不是"知道 Kafka 不好就不用",而是认得出同一系统里哪条流与日志模型对齐、哪条不对齐。订单主干(高吞吐、多组、要重放、按 key 有序)是 Kafka 的主场;通知里的复杂路由那一小块是 RabbitMQ 的主场。把它们硬塞进一个工具,就会持续撞上04 章那张表里的失败模式。

5.4自己实现(先别看参考实现)

把 §5.2 的五个决策落成代码。先合上参考实现,自己写一版——判别题的价值在于你做了选择、并能说出为什么;直接看答案等于把这章读成了配置清单。

任务

为订单系统写出两段东西:

  • 可靠 producer 配置:订单服务发 orders 事件,对账与库存不能丢、同单有序。写出 producer 端 + topic/broker 端的关键配置,并标注 key 怎么选。
  • 幂等库存消费者:消费 orders、对支付事件扣减外部库存 DB,做到重启不重复扣。写出 poll 循环 + offset 提交时机 + 去重逻辑的位置(去重落在哪一侧)。

验收 checklist(用可观测行为验,不靠"看起来对")

  • 重启不重复扣:消费者处理到一半 kill -9,重启后那批被重投,但库存 DB 的最终扣减量与只处理一次相同(幂等键挡住重复)。
  • kill 一个 broker 不丢已确认订单:producer 回调已拿到 offset 的订单,在 kill 掉一个 broker(含 leader 切换)后,下游仍能读到——没有"已 ack 却消失"的记录。
  • 同一订单严格有序:对 order-X 连发 创建→支付→发货,对账消费组读到的顺序必然是创建在前、发货在后,永不倒置。
  • 三组互不影响:把通知组 offset 重置到 24 小时前重放,库存组和对账组的进度不受任何影响。
验收陷阱

"重启不重复扣"不能靠"消费者没崩过"来证明——要主动 kill -9 制造重投再看 DB 结果。一个只在happy path 跑通的实现,会在第一次再平衡或崩溃时暴露重复扣减。验收的是失败路径的可观测行为,不是正常路径。

参考实现(写完自己版本再展开)

① 可靠 producer + topic 配置。四件套(图 4.1 的四道防线)+ 启用幂等 producer + orderId 做 key:

OrderProducer.java Java
Properties p = new Properties();
p.put("bootstrap.servers", "broker1:9092,broker2:9092,broker3:9092");
p.put("key.serializer",   "org.apache.kafka.common.serialization.StringSerializer");
p.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");

// —— 持久性四件套里的 producer 端两项(其余两项在 topic/broker 端,见下)——
p.put("acks", "all");              // 等所有 ISR 成员确认,而非 leader 单方面(堵陷阱 1)
p.put("enable.idempotence", "true"); // PID+序列号去重并强制保序;in-flight<=5 仍有序(堵正确性类乱序)
// enable.idempotence=true 会隐式要求 acks=all、retries>0,4.0 起默认即开

KafkaProducer<String, String> producer = new KafkaProducer<>(p);

void publish(OrderEvent e) {
    // 关键:key = orderId(高基数、避免热分区;同单事件同分区 → 分区内有序)
    var record = new ProducerRecord<>("orders", e.getOrderId(), e.toJson());
    producer.send(record, (md, ex) -> {
        if (ex != null) handleSendFailure(e, ex);   // 失败要可观测,别吞异常
    });
}
orders-topic.properties(topic / broker 端) Properties
# 持久性四件套的另外两项 + 不可逆截断的闸门
replication.factor=3                  # 三副本,单 broker 挂不丢分区(堵陷阱 4)
min.insync.replicas=2                 # = RF-1:挂一台仍可写、挂两台拒写(堵陷阱 3 的静默退化)
unclean.leader.election.enable=false  # 禁止落后副本上位截断已提交记录(堵陷阱 2)

# 分区数:按峰值吞吐与消费并行预估,且预留余量——加分区会破坏 key 顺序、且不可减
num.partitions=3

为什么不选事务/EOS 做 producer 端:对账与库存"不丢 + 不重复"的需求,由"四件套保证不丢" + "消费侧幂等保证不重复"组合即可满足。订单服务是纯生产者,不在"消费-转换-生产"闭环里,事务的原子 offset 提交用不上;强上事务只增延迟。

② 幂等库存消费者。关掉 auto-commit、处理后再提交、去重落在外部 DB 侧(因为副作用在 DB,Kafka 事务管不到):

InventoryConsumer.java Java
Properties c = new Properties();
c.put("bootstrap.servers", "broker1:9092,broker2:9092,broker3:9092");
c.put("group.id", "inventory");        // 独立消费组:与通知/对账各自维护 offset
c.put("key.deserializer",   "org.apache.kafka.common.serialization.StringDeserializer");
c.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
c.put("enable.auto.commit", "false");  // 关掉定时提交:否则提交的是 poll 返回的、不是处理完的(堵陷阱 5)

KafkaConsumer<String, String> consumer = new KafkaConsumer<>(c);
consumer.subscribe(List.of("orders"));

while (running) {
    var records = consumer.poll(Duration.ofMillis(500));
    for (var r : records) {
        OrderEvent e = OrderEvent.parse(r.value());
        if (e.getType() == PAID) {
            // 去重落在 DB 侧:幂等键 = orderId + 事件类型,唯一约束/upsert 挡住重投
            // 即使这批被重投(at-least-once),DB 的扣减结果与处理一次相同
            inventoryDb.deductIdempotent(e.getOrderId(), e.getSku(), e.getQty());
        }
    }
    consumer.commitSync();   // 整批处理完再提交:崩溃只会重投(可被幂等吸收),不会丢
}

每个决策"为什么选 X 不选 Y"小结:

  • key 选 orderId 不选 userId:两者都能让同单有序,但 userId 基数低、大客户订单全挤一个分区 → 热分区。orderId 高基数、分布均匀。
  • 持久性选四件套不选 acks=all 单独:单独 acks=all 在 ISR 缩到 1 时退化成 acks=1;持久性真正来自 min.insync.replicas,缺它即有丢失窗口。
  • 库存选 at-least-once+幂等不选事务 EOS:副作用在外部 DB,Kafka 事务覆盖不到外部写;DB 侧幂等键更简单且真正挡住重复。
  • 顺序选分区内有序不选全局单分区:业务只需同单有序,orderId key 已满足;单分区会把吞吐砍到 1 个分区的能力,为不需要的保证付全部代价。
  • offset 处理后提交不在处理前:先提交后处理会把失败模式从"重复"翻成"丢失";库存能容忍重复(幂等兜底)、不能容忍丢失。
亲手画一张图

合上教程,在纸上或 Excalidraw 里画出订单系统的 Kafka 架构——只画 orders topic、它的 3 个分区、以及库存/通知/对账三个消费组就行。画完回到 §5.1 的图 5.1 对照——你画的图里,给每个消费组都标了它各自的 offset 游标吗?还是把三个组画成"共享一个进度"?后者正是"把 Kafka 当队列"的典型误画:日志模型里 offset 在每个消费组各自手上,三组互不影响。

5.5反思问题

  1. 这五个决策里,哪个你做得最不确定?回看是哪一章帮你定下来的——是 01 的概念、02 的机制、还是 04 的失败模式?把"卡住→回看哪章→怎么定的"这条路径写下来。
  2. 对账组要求"不丢 + 按订单有序",通知组只要求"尽量别丢、重复无害"。如果用同一套最严格配置喂所有组,会付出什么不必要的代价?反过来,用最松的配置喂所有组,哪一组会先出事?
  3. (场景变形) 如果产品要求订单事件支持优先级——VIP 订单要插队优先处理——你上面哪些决策要改?Kafka 能优雅支持优先级吗?
反思参考(先自己想完再展开)
  1. 没有标准答案——重点是能定位到具体章节。常见的"最难"是决策 ③(消费语义):它要同时调动 02 §2.5 的 EOS 边界和 04 §4.2 的重复/乱序,并认清"副作用在外部 DB → 事务管不到"。能讲清这条回看路径,就说明判别是推理出来的、不是背的。
  2. 用最严格配置喂所有组:通知组也被迫 acks=all+RF=3+处理后提交+严格幂等,付出延迟和复杂度去保护一份"重复无害"的数据——浪费。用最松配置(acks=1、auto-commit)喂所有组:对账组先出事——它最不能丢,而 acks=1 有丢失窗口、auto-commit 会丢未处理段。结论:可靠性配置应按组的诉求分级,不是全系统一刀切。
  3. 要改决策 ④甚至整个选型。Kafka 不擅长优先级:分区日志严格按追加顺序,没有消息优先级概念,VIP 事件无法"插队"到已写入记录之前。变通办法(拆 orders-vip / orders-normal 两个 topic、消费端优先轮询 VIP)很笨重且破坏统一顺序。这正呼应 04 章表 4.1"优先级队列 → 用支持 priority 的队列系统"——优先级是把订单某条流推向 RabbitMQ 的典型信号。
进阶挑战 · 刚好够不着

给"通知组要 per-message 重试 + DLQ"设计落地方案

通知组发短信会偶发失败,要求:失败的单条消息自动重试 3 次、仍失败进死信队列人工处理,且不能阻塞同分区后面的消息。在纯 Kafka 上怎么做?这个需求是不是在提示通知组该换工具?写出两种方案并比较。

提示(卡住再展开)

纯 Kafka 方案:建 notify-retry(带延迟重试)+ notify-dlq 两个 topic,消费失败的消息发到 retry topic、计数满进 dlq——但 Kafka 没有原生 per-message 重试/DLQ,且单分区内"阻塞后面消息"很难绕(一条卡住整批)。这正是 §5.3 决策树 Q2 命中"是"的场景:通知这一段换 RabbitMQ(原生 DLQ + 单条 ack/nack + 不阻塞队列其余消息)往往比在 Kafka 上自建重试基建更省。比较维度:实现复杂度、对其余消息的阻塞、运维成本。