前三篇搭好了骨架:Kafka 是一份分布式可重放日志、Producer 怎么把消息塞进去、Consumer 怎么读出来。但有一个问题在前三篇里反复出现、却始终没讲透:
Kafka 到底怎么保证”绝不丢”?“精确一次(Exactly-Once)“又是怎么做到的、它的边界在哪?
对广告的计费、对账链路,这不是学术问题——丢一条曝光就少算一笔钱,重算一条转化就多扣一笔预算。这篇把可靠性从”配几个参数”讲到”背后的机制与边界”。
TL;DR
- “不丢”是一份三方合约:
Producer 的 acks×Broker 的复制(ISR / min.insync.replicas)×Consumer 的提交时机。任何一方松了,整条链路就不成立——只调acks=all远远不够。 acks三档:0(发了不管,最快最易丢)→1(Leader 落盘即确认,换主可能丢)→all(ISR 全复制才确认,最稳)。acks=all也可能丢:如果 ISR 收缩到只剩 Leader 一个,acks=all仍会”确认”,此时 Leader 一挂就丢。补丁是min.insync.replicas:ISR 小于它就拒绝写入,宁可拒写也不给”假成功”的确认。- 黄金组合:
acks=all+min.insync.replicas=2+RF=3——允许挂 1 个副本仍能写、确认过的数据至少 2 份、换主不丢。 unclean.leader.election.enable=false:不让”落后的、不在 ISR 里的”副本当选 Leader——这是一致性对可用性划的最后一道闸。- 幂等生产者(
enable.idempotence):给每条消息编(PID, 分区, 序列号),Broker 按序号去重、拒乱序 → 重试不重复、不乱序(默认已开)。 - 事务(
transactional.id):把”消费-处理-生产”和”提交消费位点”放进一个原子事务,靠事务协调器做两阶段提交 +read_committed隔离,实现 Kafka→Kafka 的 Exactly-Once。 - EOS 是”Kafka 内部”语义:一旦处理有外部副作用(写库、扣款、调下游),端到端”精确一次”仍需消费端幂等兜底。生产主流是 at-least-once + 幂等,而非硬上事务。
Table of contents
Open Table of contents
1. “不丢”是一份三方合约
很多人以为”消息不丢”就是把 acks 设成 all。错。不丢是一条完整的责任链,任何一环断了都会丢:
| 环节 | 谁负责 | 断了会怎样 | 关键旋钮 |
|---|---|---|---|
| 写入确认 | Producer | 不等确认就发(acks=0)→ 网络/broker 故障静默丢 | acks、retries |
| 持久与复制 | Broker | 只写了 1 份,副本没跟上就换主 → 丢 | replication.factor、min.insync.replicas、unclean.leader.election |
| 消费确认 | Consumer | 处理前就提交位点 → 崩溃丢 | enable.auto.commit、提交时机(见消费者篇) |
这篇聚焦前两环(Producer + Broker),第三环在消费者篇 §2 已讲透。记住这条链:“不丢”由三方共同保证,缺一不可。
2. acks:可靠性的总开关
acks 决定 Producer 等到什么程度才算”写成功”。
acks | 含义 | 可靠性 | 延迟 | 场景 |
|---|---|---|---|---|
0 | 发出即认为成功,不等任何确认 | 最低(会静默丢) | 最低 | 可容忍丢失的指标/日志采样 |
1 | Leader 写入本地日志即确认(不等 follower) | 中(换主窗口可能丢) | 中 | 一般吞吐优先场景 |
all(=-1) | 等 ISR 全部复制完才确认 | 最高 | 略高 | 计费/对账等不能丢的链路 |
自 3.0 起 acks 默认就是 all(且默认开幂等)。但——
只设
acks=all还不够。acks=all的语义是”等 当前 ISR 全部复制”。如果此刻 ISR 已经收缩到只剩 Leader 一个,那”全部复制”就等于”只有 Leader 有”——它一挂,确认过的消息照样丢。这就引出min.insync.replicas。
3. ISR + min.insync.replicas:补上最后一环
回顾核心原理篇 §3.7:ISR(In-Sync Replicas) 是与 Leader 保持同步的副本集合,Leader 挂了只从 ISR 里选新主。但 ISR 是动态收缩的——follower 落后太多(超 replica.lag.time.max.ms)会被踢出 ISR。
min.insync.replicas(topic/broker 级)规定:ISR 中的副本数少于这个值时,acks=all 的写入直接被拒绝(NotEnoughReplicas)。
它的哲学很硬核:宁可拒绝写入、让 Producer 报错重试或降级,也绝不给出一个”看似成功、实则随时会丢”的确认。 把”要不要冒丢失风险继续写”这个决定,从”broker 偷偷替你做”变成”你显式配置”。
黄金组合(几乎所有”不丢”场景的标配):
# topic / broker
replication.factor=3 # 3 副本
min.insync.replicas=2 # ISR 至少 2 个才接受写入
# producer
acks=all # 等 ISR 全部复制
enable.idempotence=true # 幂等(见 §5)
这套组合的含义:允许任意 1 个副本宕机,集群仍能写入;而每个被确认的消息至少存在于 2 个副本上——即使 Leader 立刻宕机,从 ISR 选出的新 Leader 也一定有这条数据。 这就是”确认即不丢”的完整保证。
为什么 RF=3 + min.insync=2,而不是 RF=2 + min.insync=2? 后者不留冗余:任何一个副本挂掉,ISR 就 < 2,topic 立刻不可写。RF=3 留了 1 个故障余量,是可用性与可靠性的经典平衡点。
4. Unclean Leader Election:一致性的最后一道闸
设想极端情况:ISR 里的副本全挂了,只剩一个不在 ISR、数据落后的 follower 还活着。要不要让它当 Leader?
unclean.leader.election.enable=false(默认,推荐):不选。分区宁可不可用,也不让落后副本当主——因为选它当主意味着丢掉它没跟上的那段已确认数据。这是”一致性优先”。unclean.leader.election.enable=true:选。牺牲一致性(丢数据)换可用性(分区继续服务)。
这是 Kafka 在 CAP 里给你的最后一个显式旋钮:要么容忍一段不可用(保住不丢),要么容忍一次数据丢失(保住可用)。 计费链路一律
false——丢钱比停一会儿严重得多。
5. 幂等生产者:重试为什么不会重复、不会乱序
acks=all + retries>0 保证了不丢,却带来一个新问题:重试会不会造成重复? 比如 Producer 发了消息、broker 也写成功了,但返回的 ack 在网络里丢了——Producer 以为失败,重发一遍 → 重复。
Kafka 3.0+(KIP-679)默认开启的幂等生产者(enable.idempotence=true) 解决了它:
机制其实朴素:
- Producer 启动时从 broker 拿一个唯一的 PID(Producer ID)。
- 发往每个分区的消息带一个单调递增的序列号(sequence number)。
- Broker 为每个
(PID, 分区)记住已接受的最大序列号:- 收到的序列号 = 已记录值(重发)→ 判重复,不追加、直接回 ack。
- 序列号跳跃/乱序 → 拒绝(
OutOfOrderSequenceException),强制按序。
幂等的边界:它只在单个 Producer 会话、单分区内保证不重不乱。Producer 重启后 PID 会变,跨会话的重复它管不了——那是事务要解决的。要幂等下保序,还需
max.in.flight.requests.per.connection ≤ 5(见 Producer 篇 §7.4)。
6. 事务:把”消费-处理-生产”变成一个原子操作
幂等解决了”单分区、单会话不重”。但典型的流处理是 consume-transform-produce(读一个 topic、处理、写另一个 topic),它涉及多个操作要么全成、要么全不成:
- 从
inputtopic 读消息; - 处理,产出结果写
outputtopic; - 提交
input的消费位点。
如果第 2 步成功、第 3 步崩了,重启会重读、重算、重复写 output → 结果重复。事务(transactional.id) 把这三件事绑成一个原子单元:
几个关键角色:
transactional.id:Producer 的稳定事务身份。跨重启保持不变,用来隔离”僵尸实例”(旧实例的写入会因 epoch 过期被拒)。- Transaction Coordinator(事务协调器):broker 上的组件,管理事务状态,状态存于内部 topic
__transaction_state。 - 两阶段提交 + 事务标记(marker):commit 时协调器先写
PREPARE_COMMIT,再向所有涉及的分区写COMMIT标记,让这批数据”对读者可见”。 isolation.level=read_committed:消费者只读已提交事务的数据;未提交/已中止的数据对它不可见(对比默认read_uncommitted)。
代码骨架:
props.put(TRANSACTIONAL_ID_CONFIG, "billing-etl-1");
props.put(ENABLE_IDEMPOTENCE_CONFIG, true); // 事务隐含要求幂等
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.initTransactions();
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
producer.beginTransaction();
try {
for (ConsumerRecord<String, String> r : records) {
producer.send(new ProducerRecord<>("billing-out", transform(r.value())));
}
// 把消费位点也纳入同一事务:处理结果与位点一起提交,原子
producer.sendOffsetsToTransaction(currentOffsets(records), consumer.groupMetadata());
producer.commitTransaction(); // 全成
} catch (Exception e) {
producer.abortTransaction(); // 全不成,回滚
}
}
7. Exactly-Once:它精确的到底是什么
把幂等 + 事务合起来,就是 Kafka 的 Exactly-Once Semantics(EOS)。Kafka Streams 一行配置 processing.guarantee=exactly_once_v2 就能开,因为它的处理天然是”读 Kafka → 算 → 写 Kafka”的闭环。
但必须把边界钉死:
EOS 精确的是”Kafka→Kafka”这段闭环:从输入 topic 到输出 topic、连同消费位点,不重不丢。它不是、也不可能是”端到端 Exactly-Once”。
原因很简单:一旦你的处理有外部副作用——写数据库、调下游 API、扣款——那个副作用发生在 Kafka 事务之外,Kafka 管不着。事务能保证”要么都写进 Kafka、要么都不写”,但没法回滚你已经发出的一笔扣款请求。
8. 端到端不重:还是得靠消费端幂等
所以生产里的现实是(这一点核心原理篇 §9 与消费者篇 §5 都强调过,这里给出完整理由):
- 纯 Kafka 内的流处理(如实时特征加工、topic 到 topic 的 ETL)→ 可以用 EOS,省心。
- 有外部副作用的消费(计费落库、扣预算、调风控)→ 主流是 at-least-once + 消费端幂等,而不是硬上事务:
- 允许重复投递;
- 用业务唯一键(
event_id)在外部存储去重:数据库唯一约束、INSERT ... ON CONFLICT、RedisSETNX、去重表。
为什么不都用事务? 事务有成本(额外的协调、
read_committed带来的可见性延迟、吞吐下降),而且跨系统它根本保证不了。用”允许重复 + 幂等消化”这套更简单、更健壮、适用面更广——这是工程上的务实选择,不是能力不足。
9. 可靠性配置速查(全链路)
把三方合约的旋钮一次列全:
| 层 | 参数 | 不丢推荐值 | 作用 |
|---|---|---|---|
| Producer | acks | all | 等 ISR 全部复制才确认 |
| Producer | enable.idempotence | true | 重试不重复、不乱序 |
| Producer | retries | 大值 / Integer.MAX_VALUE | 可重试错误自动重发 |
| Producer | max.in.flight.requests.per.connection | ≤ 5 | 幂等下保序上限 |
| Producer | delivery.timeout.ms | 按 SLA | 发送总超时(含重试) |
| Topic/Broker | replication.factor | 3 | 副本数 |
| Topic/Broker | min.insync.replicas | 2 | ISR 不足则拒写 |
| Broker | unclean.leader.election.enable | false | 不让落后副本当主 |
| Consumer | enable.auto.commit | false | 关自动提交,处理后手动提交 |
| Consumer | 提交时机 | 先处理后提交 | at-least-once(见消费者篇) |
| Consumer | isolation.level | read_committed(读事务数据时) | 只读已提交事务 |
10. 落到 AdTech:计费链路怎么配”绝不丢”
广告计费是”绝不丢”要求最高的链路——丢一条曝光 = 少收一笔钱,重算一条转化 = 多扣一笔预算。典型配置:
- 写入侧:
acks=all+enable.idempotence=true+RF=3+min.insync.replicas=2+unclean.leader.election=false。宁可高峰期偶发写入变慢/短暂拒写(Producer 重试/降级到本地缓冲),也绝不接受静默丢失。 - 消费侧:关自动提交、先落库后提交、用
event_id幂等去重(消费者篇 §5)。出问题时重放历史事件重算(核心原理篇 §6.5)——这正是”可重放日志 + 幂等”的杀手级组合:重算不怕重复,因为幂等把重复吃掉了。 - 要不要上事务? 如果是”事件流 → 特征/样本 topic”的纯 Kafka 加工,可以上 EOS;如果是”事件 → 计费落库”的外部副作用,就走 at-least-once + 幂等,别为 EOS 增加不必要的复杂度。
一句话:计费链路的可靠性 = 写入侧黄金组合(不丢)+ 消费侧幂等(不重)+ 可重放(能纠错)。 三者缺一不可。
11. 生产反模式与踩坑
- 只配
acks=all,不配min.insync.replicas:ISR 缩到 1 时仍确认,单点一挂就丢。对策:min.insync.replicas=2+RF=3。 RF=2+min.insync.replicas=2:任一副本挂 topic 立刻不可写。对策:RF 比 min.insync 至少大 1。- 图省事开
unclean.leader.election=true:换来一次”神秘丢数据”。对策:可靠链路一律false。 - 以为开了幂等/事务就端到端不重:外部副作用仍会重复。对策:消费端业务幂等(§8)。
- 事务用完不 commit/abort:挂起的事务会阻塞
read_committed消费者(LSO 卡住),LAG 飙升。对策:确保每条路径都 commit 或 abort。 - 多实例复用同一个
transactional.id:会互相把对方当”僵尸”隔离掉。对策:每个实例一个稳定且唯一的 id。 - retries=0 还想不丢:可重试错误直接失败 → 丢。对策:retries 给大值 +
delivery.timeout.ms控总时长。
12. 常见误解 ↔ 正解
| 常见误解 | 正解 |
|---|---|
acks=all 就一定不丢 | 还要 min.insync.replicas;否则 ISR 缩到 1 时仍会给”假成功”确认 |
| 副本越多越安全,RF 越大越好 | RF 越大写入越慢、成本越高;RF=3 + min.insync=2 是经典平衡点 |
| ISR 是固定的 | ISR 动态收缩/恢复,follower 落后会被移出 |
| 幂等能防一切重复 | 只保单会话、单分区不重;跨会话/跨分区靠事务 + 消费端幂等 |
| 开了事务 = 端到端 Exactly-Once | 只保 Kafka→Kafka 闭环;外部副作用仍需幂等 |
| Exactly-Once 是默认行为 | 默认是 at-least-once;EOS 要显式开幂等 + 事务 |
| unclean election 开着更安全(不会不可用) | 它用丢数据换可用;一致性链路应关闭 |
| 事务没成本,能开就开 | 有协调 / 可见性延迟 / 吞吐代价;有外部副作用时不如 at-least-once + 幂等 |
| 消息不丢只是 broker 的事 | 是 Producer × Broker × Consumer 三方合约,缺一不可 |
13. 速查表
- 不丢 = 三方合约:Producer
acks× Broker 复制(ISR/min.insync.replicas/unclean)× Consumer 提交时机。 - acks:
0最快最易丢 →1均衡 →all最稳。 - 黄金组合:
acks=all+enable.idempotence=true+RF=3+min.insync.replicas=2+unclean.leader.election=false。 min.insync.replicas:ISR 不足则拒写——宁可拒写不给假成功确认。- 幂等生产者:
(PID, 分区, seq)去重 + 拒乱序;只保单会话单分区。 - 事务:
transactional.id+ 事务协调器 + 两阶段提交 +read_committed,实现 consume-transform-produce 原子。 - EOS 边界:只保 Kafka→Kafka;有外部副作用 → at-least-once + 消费端幂等(业务唯一键)。
- 计费链路:写入侧黄金组合(不丢)+ 消费侧幂等(不重)+ 可重放(纠错)。
把可靠性讲透之后,最后一块拼图是性能——Kafka 凭什么在保证这些可靠性的同时,还能跑出百万级吞吐?下一篇Kafka 性能内核会拆开顺序写、页缓存、零拷贝与日志分段,讲清”磁盘系统为什么能有内存级速度”。
延伸阅读
本系列内部串读:
- Kafka 核心原理精讲:从一条日志到分布式流平台:ISR / High Watermark / 交付语义 / 副本机制的地基。
- Kafka Producer 深挖:分区策略与 Sticky Partitioner:
acks/enable.idempotence/max.in.flight在写入侧的取舍。 - Kafka Consumer 深挖:poll 循环、offset 提交与 rebalance:消费侧提交时机——不丢合约的第三环。
- Kafka 性能内核:磁盘系统凭什么跑出内存级吞吐:为什么持久性靠副本、而非频繁 fsync。
- 为什么 AdTech 偏爱 Kafka:广告事件中枢的五大典型场景与架构设计:“绝不丢、不算重”的三方合约在广告计费/对账链路里的落地。
一手资料与优质教程:
- Apache Kafka. Documentation — Design: Replication:ISR、Leader 选举、
min.insync.replicas、unclean election 的官方设计章节。 - Apache Kafka. KIP-98: Exactly Once Delivery and Transactional Messaging:幂等生产者与事务的一手设计提案。
- Confluent. Exactly-Once Semantics Are Possible: Here’s How Kafka Does It:EOS 的机制与边界详解。
- Apache Kafka. Producer Configs 与 Topic Configs:
acks/enable.idempotence/transactional.id/min.insync.replicas/unclean.leader.election.enable的权威说明。