Learn
Kafka/06-producer-reliability

Producer 可靠性与幂等事务

"消息会不会丢?会不会重?"是 Kafka 面试与生产环境的头号问题。答案不是一个开关,而是 acks、重试、min.insync.replicas、幂等、事务这一整条链的组合。本章把每一环的原理讲透。

1. acks:写入确认的三档语义

acks 决定 broker 在什么时机向 Producer 返回"写入成功":

acks=0:  Producer ──发──>  (不等任何确认)          最快, 网络抖动就丢
acks=1:  Producer ──发──> Leader 写入本地日志 ──确认   Leader 挂且未同步则丢
acks=all: Producer ──发──> Leader 写入
                            └─> 等 ISR 中所有副本同步完成 ──确认   最可靠
acks语义丢失场景延迟
0发出即成功网络丢包、broker 挂最低
1leader 落盘即成功leader 确认后立刻宕机,follower 未同步低
all(-1)ISR 全部同步才成功几乎不丢(配合 min.insync.replicas)较高

acks=all 单独并不安全。如果 ISR 收缩到只剩 leader 自己,acks=all 就退化成 acks=1。必须配合 topic/broker 级参数:

# ISR 中至少 2 个副本确认, 否则 Producer 收到 NotEnoughReplicas 异常
min.insync.replicas=2

黄金组合:replication.factor=3 + min.insync.replicas=2 + acks=all——允许挂 1 台 broker 且不丢数据、不停写。

2. 重试与乱序

retries=2147483647          # 默认无限重试, 由 delivery.timeout.ms 兜底
delivery.timeout.ms=120000  # 两分钟内没成功就报失败
max.in.flight.requests.per.connection=5   # 单连接未确认请求数

重试带来两个副作用:

  1. 重复:broker 已写入但 ACK 丢失,Producer 重发 → 消息重复。
  2. 乱序:max.in.flight 大于 1 时,批次 1 失败重试、批次 2 已成功 → 顺序颠倒。
乱序示意 (max.in.flight=2, 无幂等):
  发送: batch1, batch2
  batch1 失败, batch2 成功写入
  batch1 重试成功
  分区里的顺序: batch2, batch1   <- 乱了!

老式解法是 max.in.flight=1(吞吐惨烈)。现代解法是幂等 Producer,两个问题一起解决。

3. 幂等 Producer:PID + 序列号

enable.idempotence=true   # Kafka 3.x 默认已开启

原理:

  • Producer 初始化时向 broker 申请一个 PID(Producer ID)。
  • 发往每个分区的每个批次带上单调递增的序列号(sequence number)。
  • Broker 为每个"PID + 分区"记录最近序列号:
    • 收到重复序列号 → 直接丢弃(去重),仍返回成功;
    • 收到跳号 → 拒绝并报 OutOfOrderSequence,保证顺序。
Producer(PID=88) -> partition-0
  发送 seq=0 [写入]  seq=1 [写入]  seq=1 重试 [重复, 丢弃]  seq=2 [写入]
  分区日志: seq0, seq1, seq2   -- 不重不漏不乱

幂等开启后,max.in.flight 最多 5 也能保证单分区内有序。

边界:幂等只保证"单个 Producer 会话内、单分区"不重不乱。Producer 重启会拿到新 PID,重启前后的重发无法去重;跨分区、跨 topic 也不在保护范围。要更强的语义就需要事务。

4. 事务 Producer 与 exactly-once

事务提供跨分区、跨 topic 的原子写入:一批消息要么全部对消费者可见,要么全部不可见。

组件:
  transactional.id  -- 用户指定, 重启后不变, 用来找回之前的事务状态
  事务协调器 (Transaction Coordinator) -- broker 端模块
  __transaction_state -- 内部 topic, 持久化事务状态
 
流程:
  initTransactions()  向协调器注册 transactional.id, 递增 epoch
       |                (旧 epoch 的僵尸 Producer 之后的写入会被拒绝)
  beginTransaction()
  send() x N          消息正常写入各分区, 但标记为"未提交"
  commitTransaction() 协调器两阶段提交:
       |               1. 写 PREPARE_COMMIT 到 __transaction_state
       |               2. 向各分区写入 "提交标记" (control record)
  消费端 isolation.level=read_committed 只读已提交消息

Java 代码:

Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "order-svc-tx-01"); // 实例唯一且稳定
// 事务自动隐含 enable.idempotence=true, acks=all
 
Producer<String, String> producer = new KafkaProducer<>(props);
producer.initTransactions();
 
try {
    producer.beginTransaction();
    producer.send(new ProducerRecord<>("orders", "u1", "order-created"));
    producer.send(new ProducerRecord<>("order-events", "u1", "audit-log"));
    producer.commitTransaction();          // 两条消息原子可见
} catch (ProducerFencedException e) {
    producer.close();                      // 被新实例顶替(僵尸), 直接退出
} catch (KafkaException e) {
    producer.abortTransaction();           // 回滚, 消费者永远看不到这批消息
}

Python(confluent-kafka):

from confluent_kafka import Producer, KafkaException
 
p = Producer({
    "bootstrap.servers": "localhost:9092",
    "transactional.id": "order-svc-tx-01",
})
p.init_transactions()
p.begin_transaction()
try:
    p.produce("orders", key="u1", value="order-created")
    p.produce("order-events", key="u1", value="audit-log")
    p.commit_transaction()
except KafkaException:
    p.abort_transaction()
ℹ️exactly-once 的真正含义

Kafka 的 EOS(exactly-once semantics)指的是 "消费-处理-再生产" 这条 Kafka 内部链路上的精确一次(配合 sendOffsetsToTransaction,见第 20 章)。它不能让"写数据库 + 发消息"自动原子——那需要 Outbox 模式或业务侧幂等。别把 EOS 当银弹。

5. 可靠性配置矩阵

场景推荐配置
日志/埋点(允许少量丢)acks=1,retries 默认,linger 大批次
普通业务消息(不能丢,可容忍重复+下游幂等)acks=all + 幂等 + min.insync.replicas=2
金融/账务(不丢不重,链路在 Kafka 内)事务 Producer + read_committed 消费
⚠️常见坑
  1. 只设 acks=all 不设 min.insync.replicas:ISR 缩到 1 时静默退化,照样丢数据。
  2. transactional.id 多实例共用:后启动的实例会 fence 掉先启动的(ProducerFencedException),必须每实例唯一且重启保持不变(如"服务名-分区号")。
  3. 事务超时:transaction.timeout.ms(默认 60s)内没 commit 会被协调器主动 abort,长事务要调大但别超过 broker 的 transaction.max.timeout.ms。
  4. 误以为幂等能跨重启去重:重启换 PID,去重失效。跨会话不重复要靠事务或业务幂等键。

小结

  • acks=0/1/all 是延迟与可靠性的三档开关;acks=all 必须搭配 min.insync.replicas=2 才有意义
  • 黄金组合:3 副本 + min.insync.replicas=2 + acks=all
  • 幂等 Producer 用 PID + 序列号在单会话单分区内去重防乱序,3.x 默认开启
  • 事务提供跨分区原子写与 read_committed 隔离,transactional.id 要实例唯一且稳定
  • exactly-once 只覆盖 Kafka 链路内部,跨系统仍需业务幂等
🎯练习
  1. 建一个 3 副本 topic 并设 min.insync.replicas=2,停掉两个 broker 后用 acks=all 发消息,观察 NotEnoughReplicasException。
  2. 写一个事务 Producer,在两条 send 之间抛异常触发 abortTransaction,用 read_committed 与 read_uncommitted 两种消费者分别验证可见性差异。
  3. 启动两个使用相同 transactional.id 的 Producer,观察第一个实例的 ProducerFencedException。