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 下)。
没解决的:
- 不跨会话。
producerId是每次启动新申请的。生产者进程崩溃重启后会拿到新 PID,上一会话未确认的批如果重发,新 PID 无法被识别为重复。跨会话的重复要靠事务(第 05 章)用transactional.id绑定,或者靠消费端业务幂等。 - 不跨分区。序列号是按分区的,跨分区的全局顺序不存在。
- 不防丢。幂等只防重、保序,不改变
acks的丢消息条件。acks=1+ 幂等仍然会在 leader 切换时丢消息(丢的是「没复制到新 leader 的那些」)。 - 不替消费端去重。
read_committed之外的消费语义仍是至少一次,消费者重复处理一条消息是正常的,业务幂等键还是要有。
主动破坏:把 max.in.flight 调到 6
在开幂等的前提下尝试 max.in.flight.requests.per.connection=6,客户端会在启动时报配置冲突(幂等要求 ≤5)。这个练习能帮你在评审 AI 生成的「Kafka 高吞吐配置」时一眼看出矛盾——有人会把 max.in.flight 调很大来「提升吞吐」,却不知道这会和幂等打架。
生产边界
- 教学替身 vs 真实依赖:本机是单机三节点,leader 切换的触发是我手动
taskkill,不是真实的网络分区。真实环境里还有网络抖动导致的「假死」——连接没断但心跳超时,这类故障更难复现,但后果和乱序、重复是同类的。 - 要盯的指标:生产者侧的
record-error-rate、record-retry-rate、request-latency-avg,以及 broker 侧的ErrorsPerSec(含OUT_OF_ORDER_SEQUENCE_NUMBER)。retry 率突然升高通常意味着当时的 leader 不稳定,接下来就可能乱序或重复。 - 失败策略:幂等生产者遇到
OUT_OF_ORDER_SEQUENCE_NUMBER会重置 PID 并重发(可能产生重复),遇到不可重试错误会失败整批。生产端要把「发送失败」和「发送成功但可能重复」当成两种不同的异常来处理,前者重试、后者靠下游幂等。
动手
- 复现刚才的对照:对一个 1 分区 topic 分别用「幂等 + acks=all」和「非幂等 + acks=1 + max.in.flight=5 + retries=20」各发 3000 条,中途 kill leader,比较乱序点数。
- 用
dumplog找到幂等生产者写入批的producerId和baseSequence,说出这两个值在去重里各自的作用。 - 断言题:给一个
max.in.flight=10+enable.idempotence=true的配置,指出它错在哪、客户端会怎么处理。
自测
- acks=1 和 acks=all 丢消息的条件分别是什么?为什么幂等 + acks=1 仍然会丢?
- 「多个批同时在途」为什么是乱序的必要条件?把
max.in.flight设成 1 会付出什么代价? - 幂等生产者为什么能保住顺序?是从「等待」做到的,还是从「校验」做到的?
- 生产者进程崩溃重启后,同一批消息重发为什么可能重复?
transactional.id在这里补上了什么? - 你在一份 AI 生成的高吞吐配置里看到
enable.idempotence=true+max.in.flight=20。这个组合能跑起来吗?为什么?
↓ 下一步:04 章 · 消费组与再平衡