主题
概念卡片:Kafka消息可靠性不丢不重
一句话机制:消息要穿越生产 → Broker 存储 → 消费三段,每一段都可能丢;可靠性的标准答案是三段分别堵漏——生产者
acks=all + retries,Broker副本冗余 + 禁止 unclean 选举,消费者关自动提交、业务成功再手动提交;而「不重」靠手动提交引入的重复 + 消费端幂等设计兜底(Kafka 默认 at-least-once,exactly-once 需端到端配合)。
三段式丢消息排查
| 环节 | 丢消息场景 | 对策 |
|---|---|---|
| 生产者 | acks=0 不管结果;acks=1 但 Leader 挂了、Follower 没同步 | acks=all + retries 大值 + 异步回调捕获失败重发 |
| Broker 存储 | 写 PageCache 未刷盘、机器宕机 | 副本冗余(replication.factor)+ min.insync.replicas + 禁 unclean 选举 |
| 消费者 | 自动提交 offset,业务还没处理完就提交 | enable.auto.commit=false,处理成功再手动提交 |
生产者端
三种发送模式:发后即忘(fire-and-forget,效率最高可靠性最差)/ 同步(send().get() 等结果,一条等一条)/ 异步(send(msg, callback),回调里判断结果并重试)。推荐异步 + 回调 + 失败重试。
acks 三档:
| 值 | 语义 | 风险 |
|---|---|---|
| 0 | 不等服务端响应 | 网络异常/缓冲满即丢 |
| 1(默认) | Leader 写入成功即成功 | Leader 挂前 Follower 未同步 → 丢 |
| all / -1 | ISR 全部副本写入成功才成功 | ⚠️ ISR 只剩 Leader 时退化为 acks=1 |
retries 设大值 + 合理重试间隔(间隔太小网络一次波动就重试完了)。生产代码用 ListenableFuture.addCallback 打日志记录失败。
Broker 端
| 参数 | 推荐 | 说明 |
|---|---|---|
replication.factor | ≥ 3 | 保证有 Follower 可同步、可容灾 |
min.insync.replicas | > 1 | 至少写入 2 副本才算成功 |
unclean.leader.election.enable | false | 未同步副本禁止当选 Leader(0.11.0.0 起默认值由 true 改 false) |
铁律:
replication.factor > min.insync.replicas。两者相等时挂一个副本分区就不可用;生产推荐factor = min.insync.replicas + 1。
消费端
- 关自动提交(
enable.auto.commit=false),业务处理成功后再提交——否则「拿到消息但没处理完就提交」会丢消息。 - 手动提交的副作用是重复消费(处理完没提交就挂了,重启后从头再消费)→ 消费端必须幂等(唯一键去重 / 状态机判断)。
auto.offset.reset:无 offset 可提交时,earliest从分区头读(可能重复、不丢),latest从末尾读(可能丢)。- 消费状态跟踪:Kafka 用分区内一个整数 offset 标记消费进度(对比传统 MQ 的「已发送/已消费」状态标记 + 加锁维护,简单且天然支持回溯)。
事务语义与 Exactly Once
| 语义 | 说明 |
|---|---|
| at-most-once(最多一次) | 可能一次都不传 |
| at-least-once(最少一次,Kafka 默认) | 不丢但可能重复 → 业务幂等兜底 |
| exactly-once(精确一次) | 不丢不重;需将 offset 作为唯一 id 与消息处理原子化(事务性写入,如 Kafka Streams / 外部存储记录已处理 offset),成本高 |
顺序性
- 分区内有序(写入尾追加 + offset 递增 + 分区只被组内一个消费者消费);Topic 全局无序。
- 保证顺序两法:① 单分区 Topic(牺牲扩展性);② 同 key 同分区(如订单 ID 作 key,hash 取模落同一分区)。
- ⚠️ key 分区映射依赖分区数,分区数变更后同 key 可能换分区。
消费活锁(心还在、活不干)
- 现象:消费者持续发心跳但不处理消息(poll 太慢/阻塞),长期霸占分区。
- 机制:
max.poll.interval.ms活跃检测——poll 间隔超过阈值,客户端主动离组让出分区(可能抛CommitFailedException);max.poll.records限制每次 poll 返回条数。 - 对策:处理不可控时把消息处理移到独立线程,消费线程持续 poll;手动提交必须在线程处理完成后;必要时
pause()暂停分区防 OOM。
不变量(必须成立的约束)
- 可靠性三段各自堵漏,缺一环都可能丢;acks=all 不等于不丢(ISR 退化时等同 acks=1)。
- at-least-once 是 Kafka 默认语义,重复消费要靠幂等,不是靠 MQ 保证。
- factor 与 min.insync 必须留差(+1),否则高可用与可靠性不可兼得。
常见误解
- 以为
acks=all就万无一失 → ISR 只剩 Leader 时它和 acks=1 没区别。 - 以为手动提交只解决丢消息 → 它同时引入重复,必须配合幂等。
- 以为 Kafka 能保证全局有序 → 只能分区内有序。
- 以为分区越多吞吐越高 → 反噬点见 概念卡片:Kafka架构与高可用。