面试专题

Kafka 如何保证消息不丢失

从生产者、Broker 和消费者三段链路分析 Kafka 消息丢失的原因,以及 acks、ISR、副本、手动提交 Offset 和业务幂等的完整保障方案。

难度:深入更新:2026-07-23

Kafka 如何保证消息不丢失

面试回答

Kafka 保证消息不丢失需要覆盖生产者、Broker 和消费者三段链路。生产者侧使用 acks=all、开启重试和幂等生产,并检查发送回调,对最终失败进行落库、告警或补偿;Broker 侧通常设置三副本、min.insync.replicas=2,并禁止非同步副本参与 Leader 选举,保证消息在足够多的同步副本写入后才确认成功;消费者侧关闭自动提交,在业务处理成功后再提交 Offset,失败时不提交以便重新消费。这样通常实现至少一次投递,所以消费端还必须通过事件 ID、唯一索引或状态机保证幂等。涉及数据库与 Kafka 的一致性时,可以使用本地消息表或 Transactional Outbox。Kafka 的消息不丢失是端到端保障,不是只配置一个参数。

一句话总结:

生产者确认成功,Broker 多副本可靠保存,消费者处理成功后提交,业务端用幂等和补偿兜底。

详细讲解

Kafka 消息不丢失不是靠某一个参数实现的,而是一个端到端问题:

生产者必须确认消息发送成功,Broker 必须把消息可靠地保存在足够多的副本中,消费者必须在业务处理成功后再提交 Offset。

其中任何一个环节处理不当,都可能出现“业务上看起来消息丢了”的情况。

Kafka 消息不丢失的端到端保障链路

一、先明确什么叫“消息丢失”

常见的丢失场景主要有四类:

  1. 生产者调用发送方法后没有检查结果,实际发送失败却被业务忽略。
  2. Leader 收到消息后就返回成功,消息还未复制到 Follower,Leader 随后故障。
  3. 消费者先提交 Offset,再执行业务;业务处理失败后,消息不会再次投递。
  4. Kafka 中的消息没有丢,但写数据库、调用下游等业务操作失败,又没有补偿机制。

因此,Kafka 中仍能查到消息,不代表业务一定处理成功;反过来,业务没有结果,也不一定是 Broker 丢了消息。

二、生产者如何避免丢消息

1. 使用 acks=all

生产者的 acks 决定 Broker 在什么条件下确认发送成功:

配置确认时机风险
acks=0不等待 Broker 响应发送失败也无法感知,风险最高
acks=1Leader 写入本地日志后返回Leader 故障且副本尚未同步时可能丢失
acks=all当前 ISR 中的副本都确认后返回Kafka 能提供的最强生产者确认保证

关键配置:

acks=all

acks=all 并不等于绝对不丢。它还需要与副本数和 min.insync.replicas 配合,否则 ISR 中只剩 Leader 一个副本时,依然可能成功写入。

2. 开启重试和幂等生产

网络抖动、Leader 切换等临时故障可能导致发送失败,生产者需要允许重试:

enable.idempotence=true
acks=all
retries=2147483647
max.in.flight.requests.per.connection=5
delivery.timeout.ms=120000

幂等生产者会给消息附加 Producer ID 和序列号,使 Broker 能识别同一生产会话内的重复写入,避免因重试产生重复消息。

需要注意:

  • 幂等生产解决的是重试导致的重复写入,不是发送失败后的业务补偿。
  • delivery.timeout.ms 到期后,生产者仍可能最终失败。
  • 业务不能无限依赖客户端重试,最终失败必须记录、告警或进入补偿流程。

3. 必须检查发送结果

下面这种“只发送、不处理结果”的写法存在风险:

producer.send(record);

应该检查回调或 Future 的执行结果:

producer.send(record, (metadata, exception) -> {
    if (exception != null) {
        // 记录原始消息、告警并进入补偿流程
        handleSendFailure(record, exception);
    }
});

序列化失败、鉴权失败、消息过大和超时等错误,最终都需要业务明确处理。不能把“调用过 send”当成“消息已经可靠进入 Kafka”。

4. 数据库与 Kafka 的一致性

如果业务流程是:

更新数据库
→
发送 Kafka 消息

数据库更新成功、消息发送失败时,仍然会造成业务事件丢失。

常见解决方式是本地消息表,也叫 Transactional Outbox:

同一个数据库事务
├── 更新业务数据
└── 写入待发送事件表

事务提交后
→ 后台任务投递 Kafka
→ 成功后标记事件已发送

这样即使 Kafka 暂时不可用,消息也可以从本地消息表继续重试。Kafka 事务可以保证 Kafka 内部多条记录及消费 Offset 的原子性,但不会自动把外部数据库事务包含进来。

三、Broker 如何保证消息可靠保存

1. 合理设置副本数

生产环境通常使用:

replication.factor = 3

一个分区包含一个 Leader 和多个 Follower。生产者与 Leader 交互,Follower 从 Leader 复制日志。某个 Broker 故障后,可以从仍然同步的副本中选举新 Leader。

副本数提高的是容错能力,但副本只有真正保持同步才有意义。

2. 理解 ISR

ISR 是 In-Sync Replicas,即当前与 Leader 保持同步的副本集合。

例如一个三副本分区:

Leader A
Follower B
Follower C

ISR = [A, B, C]

如果 C 长时间跟不上 Leader,它会被移出 ISR:

ISR = [A, B]

acks=all 等待的是当前 ISR 的确认,而不是永远等待所有配置副本。

3. 配置 min.insync.replicas

推荐组合:

replication.factor = 3
min.insync.replicas = 2
acks = all

它表达的含义是:

至少要有两个同步副本可用,写入才允许成功。

如果 ISR 只剩一个副本,Broker 会拒绝写入。此时系统牺牲部分可用性,避免在单副本状态下继续写入并承担更高的数据丢失风险。

4. 禁止非同步副本强行成为 Leader

unclean.leader.election.enable=false

如果允许非 ISR 副本成为 Leader,在所有同步副本不可用时,Kafka 可以更快恢复分区服务,但这个副本可能缺少最新消息,从而造成数据丢失。

这是一个典型的取舍:

允许非同步副本选举:可用性更高,可能丢数据
禁止非同步副本选举:一致性更强,分区可能暂时不可用

5. acks=all 不等于每条消息都立即 fsync

Kafka 的持久性主要依赖顺序写日志、操作系统页缓存和多副本机制。Broker 返回成功,并不表示每个副本都对该消息单独执行了一次物理磁盘 fsync

实际生产中,通常通过跨 Broker 副本降低单机和单盘故障风险,而不是强制每条消息同步刷盘。若多个副本同时发生不可恢复故障,仍不存在数学意义上的绝对零丢失。

四、消费者如何避免“消费丢失”

1. 关闭自动提交 Offset

enable.auto.commit=false

自动提交可能出现以下顺序:

拉取消息
→ 自动提交 Offset
→ 执行业务
→ 进程崩溃

重启后,消费者会从已提交的下一个 Offset 继续消费,刚才尚未处理成功的消息就被跳过了。

2. 业务成功后再提交 Offset

更可靠的顺序是:

拉取消息
→ 执行业务
→ 业务成功
→ 提交 Offset

如果业务处理失败,就不提交 Offset,让消息后续重新消费:

while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(1));

    for (ConsumerRecord<String, String> record : records) {
        process(record);
    }

    consumer.commitSync();
}

示例表达的是基本原则。实际批量消费时,如果一批消息中只有部分成功,需要准确管理各分区可提交的 Offset,不能直接把失败消息之后的位置一起提交。

3. 为什么还需要业务幂等

“处理成功后提交”可以避免消息被跳过,但会引入重复:

业务处理成功
→ 提交 Offset 前进程崩溃
→ 重启后再次消费同一条消息

因此,可靠消费通常采用至少一次投递,并在业务层保证幂等。常见方式包括:

  • 使用事件 ID 建立数据库唯一索引。
  • 建立消费去重表或 Inbox 表。
  • 使用业务状态机,只允许合法状态转换。
  • 更新时带版本号或业务条件。
  • 对扣款、发券等操作使用唯一业务流水号。

不丢消息通常意味着允许重复,再通过幂等消除重复影响。

4. Kafka 内部链路的 Exactly Once

如果消费 Kafka 后仍然写回 Kafka,可以使用 Kafka 事务,把输出记录和消费 Offset 放进同一个事务:

消费消息
→ 处理
→ 写入下游 Topic
→ sendOffsetsToTransaction
→ 提交事务

下游消费者设置:

isolation.level=read_committed

这样只读取已提交事务的数据。

但如果最终写入的是 MySQL 等外部系统,Kafka 事务并不能自动保证 Kafka 与数据库的原子性,仍需要数据库事务、幂等、Outbox 或 Inbox 等方案。

五、三种投递语义

语义Offset 与业务处理顺序结果
At Most Once先提交,再处理可能丢失,通常不重复
At Least Once先处理,再提交不轻易丢失,可能重复
Exactly Once事务或幂等机制协调业务效果恰好一次

绝大多数业务系统采用:

Kafka 至少一次投递 + 消费端业务幂等。

这通常比追求所有环节的强事务更容易实现,也更便于扩展。

六、常见故障与保障措施

故障场景可能后果主要保障
生产者网络抖动发送失败或结果未知重试、幂等生产、回调处理
Leader 写入后立即故障未同步消息丢失acks=all、副本、ISR
ISR 只剩一个副本单点故障风险升高min.insync.replicas 拒绝写入
同步副本全部不可用错误选主后丢数据禁止 unclean leader election
消费者先提交后处理业务消息被跳过关闭自动提交,成功后提交
处理成功但提交前崩溃重复消费业务幂等、唯一流水号
数据库成功、Kafka 失败业务事件未发出Transactional Outbox
Kafka 成功、下游失败业务结果缺失重试、死信、补偿、告警

七、常见误区

误区一:配置 acks=all 就绝对不会丢

acks=all 只解决生产者到 Broker 的确认强度,还要配合副本数、min.insync.replicas、正确选主和消费者提交策略。

误区二:配置无限重试就不会丢

重试仍然受 delivery.timeout.ms 等条件限制,而且权限错误、序列化错误等问题无法靠盲目重试解决。最终失败必须被业务感知和补偿。

误区三:手动提交 Offset 就是 Exactly Once

手动提交只能帮助建立“处理成功后提交”的顺序,崩溃窗口仍可能造成重复消费,因此还需要业务幂等。

误区四:Kafka 事务可以自动覆盖数据库

Kafka 事务主要协调 Kafka 内部的消息与 Offset,不会自动与 MySQL 等外部数据库形成同一个原子事务。

误区五:Broker 返回成功就等于所有副本已经物理刷盘

Kafka 的确认、副本复制和磁盘刷盘是不同概念。高可靠主要依靠多副本和同步副本约束,而不是把每条消息都单独同步刷盘。

八、监控和运维同样重要

配置正确后,还应持续监控:

  • 发送失败率、重试次数和发送延迟。
  • ISR 收缩、未充分复制分区和离线分区。
  • Broker 磁盘空间、磁盘延迟和网络异常。
  • 消费积压、消费失败、重平衡次数。
  • 死信消息和补偿任务堆积。
  • 业务事件与最终结果的对账差异。

没有监控和对账,即使消息已经丢失,系统也可能长期无法发现。

参考资料

DISCUSSION

评论与讨论

留下你的想法