KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA

03 · 生产者的确认与顺序 — keel 龙骨

生产端最容易被「配一个参数就以为解决了」的地方。acks 决定丢不丢,max.in.flight 和 retries 决定乱不乱,而幂等生产者只解决其中一部分。这一章把这三组参数的边界钉清楚,并且真的跑出一次乱序。

生产端最容易被「配一个参数就以为解决了」的地方。acks 决定丢不丢,max.in.flight 和 retries 决定乱不乱,而幂等生产者只解决其中一部分。这一章把这三组参数的边界钉清楚,并且真的跑出一次乱序。

现场

订单的同一张单会有多条状态变更(创建、支付、发货),它们用同一个 key 发到同一个分区,下游要按顺序处理。生产者在发这批消息时遭遇了一次 leader 切换(第 02 章里 a2 被停掉)。问题是:切换之后,同一个 key 的消息还能保证顺序吗?

先说结论:取决于配置。同一台机器、同一份数据、同一次故障注入,acks=1 + 重试跑出了乱序,幂等生产者跑出了严格顺序。

非幂等 + 重试:乱序是怎么产生的

flowchart TD
  P["① 生产者<br/>acks=1, max.in.flight=5, retries=20"] --> A["② 批 seq 1303 -> leader a2<br/>收到 ack"]
  A --> B["③ 批 seq 1304、1305 同时发出<br/>仍发给 a2,在途未确认"]
  B -. "a2 被 kill,两个批的 ack 随连接一起丢" .-> X["④ 连接断开"]
  X --> C["⑤ 刷新元数据<br/>新 leader 变成 a3"]
  C --> D["⑥ 重试 seq 1305 -> a3<br/>先成功"]
  D --> E["⑦ 重试 seq 1304 -> a3"]
  E --> F["⑧ 重试 seq 1303 -> a3"]
  F --> G["⑨ orders-14-0 里物理顺序<br/>...1305, 1304, 1303<br/>与发送顺序相反"]
  style P fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
  style A fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
  style B fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
  style X fill:#ffe3e3,stroke:#c92a2a,color:#1a1a2e
  style C fill:#fff9db,stroke:#f08c00,color:#1a1a2e
  style D fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
  style E fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
  style F fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
  style G fill:#ffe3e3,stroke:#c92a2a,color:#1a1a2e

关键是「多个批同时在途」:max.in.flight=5 允许 5 个批同时飞。leader 死掉时,这几个批的 ack 全丢,它们各自重试,重试落到新 leader 的顺序不再受原始发送顺序保护,于是日志里的物理顺序被打乱。

幂等生产者:批头里的序列号

flowchart TD
  P["① 生产者<br/>enable.idempotence=true, acks=all, max.in.flight=5"] --> I["② init 时向 broker 申请到 producerId=PID"]
  I --> A["③ 批(PID, baseSeq=0, lastSeq=72) -> leader a3"]
  A --> B["④ broker 记下该 PID 的 nextSeq=73<br/>返回 ack"]
  B -. "leader 切换,这个批的 ack 丢失" .-> R["⑤ 重试同一批(PID, baseSeq=0)"]
  R --> D{"⑥ broker 校验 baseSeq<br/>是否等于期望的 nextSeq?"}
  D -- "重复(序列号已存在)" --> E["⑦ 丢弃该批,返回上次的 ack"]
  D -- "等于 nextSeq" --> F["⑦ 正常写入并推进 nextSeq"]
  E --> G["⑧ 日志里没有重复记录"]
  F --> G
  G --> H["⑨ 后续批(PID, baseSeq=73...) 依次写入<br/>日志顺序严格 = 发送顺序"]
  style P fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
  style I fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
  style A fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
  style B fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
  style R fill:#fff9db,stroke:#f08c00,color:#1a1a2e
  style D fill:#fff9db,stroke:#f08c00,color:#1a1a2e
  style E fill:#f3f0ff,stroke:#7048e8,color:#1a1a2e
  style F fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
  style G fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
  style H fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e

幂等生产者做三件事:给每个生产者实例分配 producerId;给每条记录分配批内 sequence;broker 侧维护每个 PID 的「下一个期望序列号」。重复的批被识别为重复直接丢弃(返回上次的 ack),序列号不连续(跳号或逆序)的批被拒绝。去重和保序是同一套机制的两个结果。

精确定义

acks 四档:

取值 语义 丢消息的条件
0 发出即返回,不等任何确认 broker 没收到 / 落盘前宕机
1 等 leader 写入本地日志就返回 leader 落盘后、follower 复制前宕机,且新 leader 没有这条
all(等价 -1) 等 ISR 全部 ack ISR 里最后一个持有者也丢(ISR 缩到 1 时)
all + min.insync.replicas=2 等 ISR 全部 ack,且 ISR<2 直接拒写 两个持有者同时永久丢失

注意「acks=1」那一行的坑:leader ack 之后这条消息可能还没被任何 follower 复制,此时 leader 宕机、controller 从 ISR 里另选一个还没收到这条消息的副本当 leader,这条消息就没了。本机第 06 章会看到同类丢失的现场。

max.in.flight.requests.per.connection:一个连接上未确认请求的上限。它 > 1 才有吞吐优势,也才有乱序风险。幂等开启后可以设到 5 而仍保序(Kafka 用序列号兜底),超过 5 会被客户端拒绝。

enable.idempotence:开幂等会自动要求 acks=all、retries>0、max.in.flight<=5。也就是说开幂等的同时你顺带接受了 acks=all 的可用性代价。

证据一:无故障时幂等生产者保序

orders-07(1 分区 3 副本),幂等生产者发 3000 条 seq=NNNNN(固定 key,全落同一分区):

[idem] idempotent=True acks=all max_in_flight=5 retries=10
[idem] flush ok, sent=3000
读回 3000 条
重复(同一 seq 出现多次): 0 个
乱序点(前一个 seq > 后一个 seq): 0 处
seq 范围 [0,2999], 缺失 0 个

零重复、零乱序、零缺失。

证据二:非幂等 + leader 切换 → 真实乱序

orders-14(同类 topic),生产者 enable.idempotence=false, acks=1(整数), max_in.flight=5, retries=20,摊到 12 秒发送,中途 kill 掉它的 leader:

[noidem2] topic=orders-14 leader=a2 (pid=51312) idem=0 acks=1 mif=5
[noidem2] kill a2 at 15:47:33
[noidem2] sent 3000 (+15.1s)
[noidem2] flush ok, sent=3000
读回 3000 条
重复(同一 seq 出现多次): 0 个
乱序点(前一个 seq > 后一个 seq): 2 处, 例: [(1305, 1304), (1304, 1303)]
seq 范围 [0,2999], 缺失 0 个

3000 条一条没丢、一条没重,但出现 2 处乱序,正好落在 leader 切换的那个位置(seq 1303~1305)。这正是上面那张链路图预测的结果。乱序不是「偶发的坏运气」,而是 max.in.flight>1 + 重试 + leader 切换这个组合的必然产物。

证据三:同样的故障注入,幂等生产者零乱序

orders-15,enable.idempotence=true, acks=all, max.in.flight=5, retries=20,同样在发送中途 kill leader:

[idem2] topic=orders-15 leader=a3 (pid=73548) idem=1 acks=all mif=5
[idem2] kill a3 at 15:48:54
[idem2] flush ok, sent=3000
读回 3000 条
重复(同一 seq 出现多次): 0 个
乱序点(前一个 seq > 后一个 seq): 0 处
seq 范围 [0,2999], 缺失 0 个

对照组和实验组唯一的差别就是那个开关。dump 幂等生产者的段,能看到批头里带着身份:

baseOffset: 0 lastOffset: 0 count: 1 baseSequence: 0 lastSequence: 0
producerId: 4000 producerEpoch: 0 ... magic: 2 ...
| offset: 0 ... sequence: 0 ... key: k payload: seq=00000

producerId: 4000 是这次生产者实例的身份,baseSequence: 0 是这一批的起始序列号——第 02 章说 broker 靠这两个字段做去重,这就是它们在磁盘上的样子。

一次踩坑记录:客户端参数类型

写这组实验时我先用 acks="1"(字符串)配 kafka-python,结果生产端静默地什么都没写:send() 不报错、flush() 也不报错,但 topic 的 LEO 一直是 0。把三种写法摆在一起对比:

acks=1      -> partition=0 offset=4
acks=1      -> partition=0 offset=5
acks='all'  -> partition=0 offset=6
acks='all'  -> partition=0 offset=7
acks='1'    -> KafkaTimeoutError: Timeout after waiting for 8 secs.
acks='1'    -> KafkaTimeoutError: Timeout after waiting for 8 secs.

acks 必须是整数或字符串 'all';传字符串 '1' 会导致批永远不被发出,只有逐条 future.get() 才能看出异常。这不是 Kafka 的语义问题,是客户端库的参数校验缺失。排查「发出去没落地」这类问题时,先 await 每一条 future 拿到 offset,别只看 flush() 是否返回。

幂等生产者解决了什么、没解决什么

解决的:单个生产者会话内,对一个分区,重复发送不会产生重复记录(broker 按 PID+sequence 去重),且写入顺序与发送顺序一致(在 max.in.flight<=5 下)。

没解决的:

主动破坏:把 max.in.flight 调到 6

在开幂等的前提下尝试 max.in.flight.requests.per.connection=6,客户端会在启动时报配置冲突(幂等要求 ≤5)。这个练习能帮你在评审 AI 生成的「Kafka 高吞吐配置」时一眼看出矛盾——有人会把 max.in.flight 调很大来「提升吞吐」,却不知道这会和幂等打架。


生产边界

动手

  1. 复现刚才的对照:对一个 1 分区 topic 分别用「幂等 + acks=all」和「非幂等 + acks=1 + max.in.flight=5 + retries=20」各发 3000 条,中途 kill leader,比较乱序点数。
  2. 用 dumplog 找到幂等生产者写入批的 producerId 和 baseSequence,说出这两个值在去重里各自的作用。
  3. 断言题:给一个 max.in.flight=10 + enable.idempotence=true 的配置,指出它错在哪、客户端会怎么处理。

自测

  1. acks=1 和 acks=all 丢消息的条件分别是什么?为什么幂等 + acks=1 仍然会丢?
  2. 「多个批同时在途」为什么是乱序的必要条件?把 max.in.flight 设成 1 会付出什么代价?
  3. 幂等生产者为什么能保住顺序?是从「等待」做到的,还是从「校验」做到的?
  4. 生产者进程崩溃重启后,同一批消息重发为什么可能重复?transactional.id 在这里补上了什么?
  5. 你在一份 AI 生成的高吞吐配置里看到 enable.idempotence=true + max.in.flight=20。这个组合能跑起来吗?为什么?

↓ 下一步:04 章 · 消费组与再平衡

进入 keel 阅读