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 挂 | 最低 |
| 1 | leader 落盘即成功 | 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 # 单连接未确认请求数重试带来两个副作用:
- 重复:broker 已写入但 ACK 丢失,Producer 重发 → 消息重复。
- 乱序:
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 消费 |
⚠️常见坑
- 只设 acks=all 不设 min.insync.replicas:ISR 缩到 1 时静默退化,照样丢数据。
- transactional.id 多实例共用:后启动的实例会 fence 掉先启动的(ProducerFencedException),必须每实例唯一且重启保持不变(如"服务名-分区号")。
- 事务超时:
transaction.timeout.ms(默认 60s)内没 commit 会被协调器主动 abort,长事务要调大但别超过 broker 的transaction.max.timeout.ms。 - 误以为幂等能跨重启去重:重启换 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 链路内部,跨系统仍需业务幂等
🎯练习
- 建一个 3 副本 topic 并设 min.insync.replicas=2,停掉两个 broker 后用 acks=all 发消息,观察 NotEnoughReplicasException。
- 写一个事务 Producer,在两条 send 之间抛异常触发 abortTransaction,用 read_committed 与 read_uncommitted 两种消费者分别验证可见性差异。
- 启动两个使用相同 transactional.id 的 Producer,观察第一个实例的 ProducerFencedException。