Kafka 如何保证消息不丢失
从生产者、Broker 和消费者三段链路分析 Kafka 消息丢失的原因,以及 acks、ISR、副本、手动提交 Offset 和业务幂等的完整保障方案。
Kafka 如何保证消息不丢失
面试回答
Kafka 保证消息不丢失需要覆盖生产者、Broker 和消费者三段链路。生产者侧使用
acks=all、开启重试和幂等生产,并检查发送回调,对最终失败进行落库、告警或补偿;Broker 侧通常设置三副本、min.insync.replicas=2,并关闭非同步副本选举,即配置unclean.leader.election.enable=false,保证消息在足够多的同步副本写入后才确认成功;消费者侧关闭自动提交,在业务处理成功后再提交 Offset,失败时不提交以便重新消费。这样通常实现至少一次投递,所以消费端还必须通过事件 ID、唯一索引或状态机保证幂等。涉及数据库与 Kafka 的一致性时,可以使用 Transactional Outbox;消费端需要应对重复消息时,可以使用 Inbox。Kafka 的消息不丢失是端到端保障,不是只配置一个参数。
一句话总结:
生产者确认成功,Broker 多副本可靠保存,消费者处理成功后提交,业务端用幂等和补偿兜底。
一张图看懂
详细讲解
Kafka 消息不丢失不是靠某一个参数实现的,而是一个端到端问题:
生产者必须确认消息发送成功,Broker 必须把消息可靠地保存在足够多的副本中,消费者必须在业务处理成功后再提交 Offset。
其中任何一个环节处理不当,都可能出现“业务上看起来消息丢了”的情况。
一、先明确什么叫“消息丢失”
常见的丢失场景主要有四类:
- 生产者调用发送方法后没有检查结果,实际发送失败却被业务忽略。
- Leader 收到消息后就返回成功,消息还未复制到 Follower,Leader 随后故障。
- 消费者先提交 Offset,再执行业务;业务处理失败后,消息不会再次投递。
- Kafka 中的消息没有丢,但写数据库、调用下游等业务操作失败,又没有补偿机制。
因此,Kafka 中仍能查到消息,不代表业务一定处理成功;反过来,业务没有结果,也不一定是 Broker 丢了消息。
二、生产者如何避免丢消息
1. 使用 acks=all
生产者的 acks 决定 Broker 在什么条件下确认发送成功:
| 配置 | 确认时机 | 风险 |
|---|---|---|
acks=0 | 不等待 Broker 响应 | 发送失败也无法感知,风险最高 |
acks=1 | Leader 写入本地日志后返回 | 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. 关闭非同步副本选举
unclean.leader.election.enable=false
unclean leader election 指允许不在 ISR 中的副本被选举为 Leader。如果所有同步副本都不可用,非 ISR 副本可能缺少最新消息;将它选为 Leader 虽然能更快恢复分区服务,却可能造成日志回退和数据丢失。
因此,对数据可靠性要求较高的场景应明确配置:
关闭非同步副本选举:
unclean.leader.election.enable=false。
这是一个典型的取舍:
开启非同步副本选举:可用性更高,但可能丢失消息
关闭非同步副本选举:数据可靠性更强,但分区可能暂时不可用
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 应对重复消费,具体区别见下一节。
五、Outbox 与 Inbox 分别是什么
Outbox 和 Inbox 都借助数据库本地事务解决消息链路中的一致性问题,但它们位于不同位置:
业务生产方 → Outbox → Kafka → Inbox → 业务消费方
1. Outbox:保证业务事件可靠发出
Outbox 位于消息生产方,用来解决“数据库更新成功,但 Kafka 消息发送失败”的双写不一致问题。
核心流程:
同一个数据库事务
├── 更新业务数据
└── 写入 Outbox 事件表
事务提交
→ 后台投递程序扫描 Outbox
→ 发送 Kafka
→ 收到成功确认后标记已发送
例如创建订单时,不直接把“订单创建事件”的可靠性寄托在一次 Kafka 调用上,而是在创建订单的同一个数据库事务中写入一条 Outbox 记录。即使 Kafka 暂时不可用,后台任务仍可继续重试。
Outbox 的关键点:
- 业务数据和事件记录必须写入同一个数据库事务。
- 每条事件应有全局唯一的
eventId。 - 投递程序可以轮询事件表,也可以通过 CDC 捕获新增事件。
- “发送成功但标记失败”可能导致再次投递,因此消费者仍需幂等。
- Outbox 解决的是事件可靠发出,不保证消费者一定处理成功。
这里的 CDC 是 Change Data Capture(变更数据捕获)。它会读取数据库的变更日志,例如 MySQL Binlog,捕获数据的新增、修改和删除,再把变更发送到 Kafka。常用工具有 Debezium、Canal。
在 Outbox 场景中,它的工作链路是:
业务事务写入业务表和 Outbox 表
→ CDC 监听 Binlog
→ 捕获新增的 Outbox 记录
→ 发送到 Kafka
相比定时扫描 Outbox 表,CDC 延迟更低,也能减少频繁查询数据库的压力。但消息链路仍可能产生重复,因此必须保留 eventId,并由消费端保证幂等。
2. Inbox:保证重复消息只产生一次业务效果
Inbox 位于消息消费方,用来记录已经处理过的事件,主要解决至少一次投递带来的重复消费问题。
核心流程:
收到 Kafka 消息
→ 开启数据库事务
→ 根据 eventId 写入 Inbox 记录
→ 执行业务更新
→ 提交数据库事务
→ 提交 Kafka Offset
Inbox 表通常对 eventId 建立唯一索引。重复消息到达时,如果发现该事件已经存在,就不再重复执行业务逻辑。
Inbox 的关键点:
- Inbox 记录和业务更新必须处于同一个数据库事务。
- 唯一索引负责防止同一事件被并发重复处理。
- 数据库事务提交后、Offset 提交前发生故障,消息仍会重投,但 Inbox 可以识别重复。
- 如果消费逻辑还会调用外部 HTTP 接口,Inbox 无法自动让数据库和外部接口形成原子事务;仍需下游幂等、重试或补偿。
3. 两者的区别
| 对比项 | Outbox | Inbox |
|---|---|---|
| 所在位置 | 消息生产方 | 消息消费方 |
| 主要目标 | 防止业务成功但事件未发出 | 防止重复消息重复产生业务效果 |
| 本地事务内容 | 业务更新 + 写事件表 | 写消费记录 + 业务更新 |
| 后台动作 | 把待发送事件投递到 Kafka | 通常不需要再次投递 |
| 仍需注意 | 投递过程可能产生重复消息 | 外部调用仍需幂等或补偿 |
一句话区分:
Outbox 保证“该发的消息最终发出去”,Inbox 保证“重复收到的消息不会重复生效”。
六、三种投递语义
| 语义 | Offset 与业务处理顺序 | 结果 |
|---|---|---|
| At Most Once | 先提交,再处理 | 可能丢失,通常不重复 |
| At Least Once | 先处理,再提交 | 不轻易丢失,可能重复 |
| Exactly Once | 事务或幂等机制协调 | 业务效果恰好一次 |
绝大多数业务系统采用:
Kafka 至少一次投递 + 消费端业务幂等。
这通常比追求所有环节的强事务更容易实现,也更便于扩展。
七、常见故障与保障措施
| 故障场景 | 可能后果 | 主要保障 |
|---|---|---|
| 生产者网络抖动 | 发送失败或结果未知 | 重试、幂等生产、回调处理 |
| Leader 写入后立即故障 | 未同步消息丢失 | acks=all、副本、ISR |
| ISR 只剩一个副本 | 单点故障风险升高 | min.insync.replicas 拒绝写入 |
| 同步副本全部不可用 | 非同步副本选主后日志回退 | 关闭非同步副本选举:unclean.leader.election.enable=false |
| 消费者先提交后处理 | 业务消息被跳过 | 关闭自动提交,成功后提交 |
| 处理成功但提交前崩溃 | 重复消费 | 业务幂等、唯一流水号 |
| 数据库成功、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 磁盘空间、磁盘延迟和网络异常。
- 消费积压、消费失败、重平衡次数。
- 死信消息和补偿任务堆积。
- 业务事件与最终结果的对账差异。
没有监控和对账,即使消息已经丢失,系统也可能长期无法发现。
评论与讨论
回复 :
留下你的想法