Skip to content
Charles Shao
Go back

Kafka 可靠性与 Exactly-Once 深挖:acks、ISR、min.insync.replicas 与事务,把『不丢不重』讲到底

views

前三篇搭好了骨架:Kafka 是一份分布式可重放日志Producer 怎么把消息塞进去Consumer 怎么读出来。但有一个问题在前三篇里反复出现、却始终没讲透:

Kafka 到底怎么保证”绝不丢”?“精确一次(Exactly-Once)“又是怎么做到的、它的边界在哪?

对广告的计费、对账链路,这不是学术问题——丢一条曝光就少算一笔钱,重算一条转化就多扣一笔预算。这篇把可靠性从”配几个参数”讲到”背后的机制与边界”。

TL;DR

Table of contents

Open Table of contents

1. “不丢”是一份三方合约

很多人以为”消息不丢”就是把 acks 设成 all。错。不丢是一条完整的责任链,任何一环断了都会丢:

环节谁负责断了会怎样关键旋钮
写入确认Producer不等确认就发(acks=0)→ 网络/broker 故障静默丢acksretries
持久与复制Broker只写了 1 份,副本没跟上就换主 → 丢replication.factormin.insync.replicasunclean.leader.election
消费确认Consumer处理前就提交位点 → 崩溃丢enable.auto.commit、提交时机(见消费者篇

这篇聚焦前两环(Producer + Broker),第三环在消费者篇 §2 已讲透。记住这条链:“不丢”由三方共同保证,缺一不可。

2. acks:可靠性的总开关

acks 决定 Producer 等到什么程度才算”写成功”

Kafka acks 三档可靠性对比图,自上而下三块。第一块 acks=0(发了就不管,fire-and-forget):Producer 发送即返回、不等确认,注解为最快但最易丢,网络或 broker 故障时消息无声消失。第二块 acks=1(等 Leader 落盘就确认):Producer 写入 Leader,Leader 不等 follower 就返回 ack,注解为均衡选择,但 Leader 确认后若在 follower 复制前宕机仍可能丢。第三块 acks=all(等 ISR 全部复制才确认,配 min.insync.replicas):Producer 写入 Leader,Leader 复制给 ISR 内的 followers,所有 ISR 复制完、HW 推进后 Leader 才返回 ack,注解为最稳,确认即代表已被所有 ISR 持有、换主也不丢,代价是延迟略高。

acks含义可靠性延迟场景
0发出即认为成功,不等任何确认最低(会静默丢)最低可容忍丢失的指标/日志采样
1Leader 写入本地日志即确认(不等 follower)中(换主窗口可能丢)一般吞吐优先场景
all(=-1ISR 全部复制完才确认最高略高计费/对账等不能丢的链路

自 3.0 起 acks 默认就是 all(且默认开幂等)。但——

只设 acks=all 还不够。 acks=all 的语义是”等 当前 ISR 全部复制”。如果此刻 ISR 已经收缩到只剩 Leader 一个,那”全部复制”就等于”只有 Leader 有”——它一挂,确认过的消息照样丢。这就引出 min.insync.replicas

3. ISR + min.insync.replicas:补上最后一环

回顾核心原理篇 §3.7ISR(In-Sync Replicas) 是与 Leader 保持同步的副本集合,Leader 挂了只从 ISR 里选新主。但 ISR 是动态收缩的——follower 落后太多(超 replica.lag.time.max.ms)会被踢出 ISR。

ISR 与 min.insync.replicas 协作图,分上下两块。上块『ISR 健康:RF=3, ISR={L,F1,F2}, min.insync.replicas=2』:Leader 与 Follower1、Follower2 都在 ISR 内,注解为 ISR 大小 3 ≥ 2,接受写入,acks=all 确认表示至少 2 副本已持有,安全。下块『ISR 收缩到 1:两个 follower 掉队被移出 ISR』:只剩 Leader 在 ISR,Follower1、Follower2 因落后被移出,注解为 ISR 大小 1 小于 min.insync.replicas=2,于是 Producer 写入被直接拒绝(NotEnoughReplicas),宁可拒写也不给假装成功却随时会丢的确认。底部蓝色关键注解:黄金组合 acks=all + min.insync.replicas=2 + RF=3,允许挂 1 个副本仍能写、且确认过的数据至少 2 份、换主不丢;只配 acks=all 不配 min.insync.replicas,ISR 缩到 1 时仍会确认,单点一挂就丢。

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?

这是 Kafka 在 CAP 里给你的最后一个显式旋钮:要么容忍一段不可用(保住不丢),要么容忍一次数据丢失(保住可用)。 计费链路一律 false——丢钱比停一会儿严重得多。

5. 幂等生产者:重试为什么不会重复、不会乱序

acks=all + retries>0 保证了不丢,却带来一个新问题:重试会不会造成重复? 比如 Producer 发了消息、broker 也写成功了,但返回的 ack 在网络里丢了——Producer 以为失败,重发一遍 → 重复。

Kafka 3.0+(KIP-679)默认开启的幂等生产者(enable.idempotence=true 解决了它:

幂等生产者去重时序图,参与方为 Producer(PID=42)和 Broker(Partition Leader)。前提注解:每条消息带 (PID, 分区, 序列号 seq),Broker 记住每个 PID 每分区的最大 seq。流程:① Producer 写入 seq=5;② Broker 已接受并记录 last seq=5;③ Producer 写入 seq=6;④ 网络抖动导致 ack 丢失,Producer 没收到确认;⑤ Producer 重试 seq=6(同一条),Broker 发现 seq=6 等于已记录的 last seq,判定重复,不再追加直接回 ack,实现不重复;⑥ Producer 写入 seq=8(乱序,跳过了 7),Broker 发现 seq 不连续,拒绝并报 OutOfOrderSequence,强制按序,从而不乱序,Producer 再按序重发。

机制其实朴素:

幂等的边界:它只在单个 Producer 会话、单分区内保证不重不乱。Producer 重启后 PID 会变,跨会话的重复它管不了——那是事务要解决的。要幂等下保序,还需 max.in.flight.requests.per.connection ≤ 5(见 Producer 篇 §7.4)。

6. 事务:把”消费-处理-生产”变成一个原子操作

幂等解决了”单分区、单会话不重”。但典型的流处理是 consume-transform-produce(读一个 topic、处理、写另一个 topic),它涉及多个操作要么全成、要么全不成

  1. input topic 读消息;
  2. 处理,产出结果写 output topic;
  3. 提交 input 的消费位点。

如果第 2 步成功、第 3 步崩了,重启会重读、重算、重复写 output → 结果重复。事务(transactional.id 把这三件事绑成一个原子单元:

Kafka 事务流程时序图,参与方为 Producer(transactional.id)、Transaction Coordinator、数据分区(output topic)、__consumer_offsets。前提注解:把消费-处理-生产放进一个事务。流程:① Producer 调 initTransactions 拿到 producer epoch 以隔离僵尸实例;② beginTransaction;③ 向 output 数据分区写处理结果(标记为事务中、未提交);④ sendOffsetsToTransaction 把消费位点也纳入同一事务;⑤ 协调器把待提交的消费位点写入 __consumer_offsets;⑥ commitTransaction;⑦⑧ 两阶段提交:协调器先写 PREPARE_COMMIT,再向所有涉及分区(output 数据分区与 __consumer_offsets)写 COMMIT 事务标记 marker。末尾注解:只有 read_committed 的消费者能看到已提交事务的数据,未提交或已中止的数据对其不可见。

几个关键角色:

代码骨架:

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 都强调过,这里给出完整理由):

为什么不都用事务? 事务有成本(额外的协调、read_committed 带来的可见性延迟、吞吐下降),而且跨系统它根本保证不了。用”允许重复 + 幂等消化”这套更简单、更健壮、适用面更广——这是工程上的务实选择,不是能力不足。

9. 可靠性配置速查(全链路)

把三方合约的旋钮一次列全:

参数不丢推荐值作用
Produceracksall等 ISR 全部复制才确认
Producerenable.idempotencetrue重试不重复、不乱序
Producerretries大值 / Integer.MAX_VALUE可重试错误自动重发
Producermax.in.flight.requests.per.connection≤ 5幂等下保序上限
Producerdelivery.timeout.ms按 SLA发送总超时(含重试)
Topic/Brokerreplication.factor3副本数
Topic/Brokermin.insync.replicas2ISR 不足则拒写
Brokerunclean.leader.election.enablefalse不让落后副本当主
Consumerenable.auto.commitfalse关自动提交,处理后手动提交
Consumer提交时机先处理后提交at-least-once(见消费者篇
Consumerisolation.levelread_committed(读事务数据时)只读已提交事务

10. 落到 AdTech:计费链路怎么配”绝不丢”

广告计费是”绝不丢”要求最高的链路——丢一条曝光 = 少收一笔钱,重算一条转化 = 多扣一笔预算。典型配置:

一句话:计费链路的可靠性 = 写入侧黄金组合(不丢)+ 消费侧幂等(不重)+ 可重放(能纠错)。 三者缺一不可。

11. 生产反模式与踩坑

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. 速查表

把可靠性讲透之后,最后一块拼图是性能——Kafka 凭什么在保证这些可靠性的同时,还能跑出百万级吞吐?下一篇Kafka 性能内核会拆开顺序写、页缓存、零拷贝与日志分段,讲清”磁盘系统为什么能有内存级速度”。


延伸阅读

本系列内部串读:

一手资料与优质教程:


views
Share this post on:

Previous Post
为什么 AdTech 偏爱 Kafka:广告事件中枢的五大典型场景与架构设计
Next Post
Kafka 性能内核:磁盘系统凭什么跑出内存级吞吐