消息队列01:Kafka
从分区机制、ISR 副本同步、消费者组 Rebalance 到 Exactly-Once 语义,拆解 Kafka 高吞吐背后的设计取舍和常见坑点。
Kafka 不是”消息队列”这个分类里最老的产品,但它把”高吞吐日志管道”这件事做到了极致。很多项目引入 Kafka,是因为日志量太大、埋点太多、或者要把数据喂给 Flink/Spark 做实时计算。
这篇文章不列所有参数,而是围绕四个核心机制:分区、ISR、消费者组、Exactly-Once,把 Kafka 的设计逻辑和常见坑点讲清楚。
为什么是 Kafka,而不是别的
先定位清楚:Kafka 的核心场景是海量数据日志/事件流的管道,追求极致吞吐量,适合 ETL、日志收集、事件溯源、流式计算。它不是为”金融交易消息”设计的,那是 RocketMQ 的主场。
两者的取舍很直接:
| 维度 | Kafka | RocketMQ |
|---|---|---|
| 存储模型 | 分区式存储,每 Partition 独立文件 | 集中式 CommitLog + ConsumeQueue 索引 |
| Topic 数量敏感度 | Topic > 50 时 PageCache 污染严重,I/O 性能急剧下降 | 单机支持上千 Topic 无明显衰减 |
| 消息可靠性 | 依赖 ISR + acks=all,需调优 | 同步刷盘 + 同步复制,业务消息可靠性更强 |
| 事务消息 | 支持,但偏流处理,缺乏服务端回查 | 原生事务消息 + 事务回查,贴近交易场景 |
| 延迟消息 | 不原生支持,需额外实现时间轮 | 原生支持 18 个延迟级别 |
简单判断:日均数据量 TB/PB 级、日志采集、实时计算、埋点流处理 → 选 Kafka。订单、支付、交易链路 → 选 RocketMQ。两者可以共存,业务消息走 RocketMQ,日志和数据管道走 Kafka。
核心架构:Producer、Broker、Topic、Partition、Consumer Group
Kafka 的架构不复杂,但每个组件都有明确职责:
- Producer:消息生产者,决定消息发到哪个 Partition(分区策略)
- Broker:Kafka 节点,负责存储和读写,多个 Broker 组成集群
- Topic:逻辑上的消息分类,物理上拆成多个 Partition
- Partition:物理存储单元,是有序的、不可变的追加日志
- Replica:每个 Partition 可以配置副本,Leader 负责读写,Follower 做数据同步
- Consumer Group:消费者组,组内每个 Partition 只被一个 Consumer 消费,实现并行扩展
核心设计思路:分区做水平扩展,副本做高可用,顺序追加做高吞吐。
Partition 数量决定了消费者的最大并行度,消息顺序仅保证在 Partition 内有序。如果业务要求全局有序,只能单 Partition 单 Consumer,吞吐会很低。大多数场景只需要”局部有序”,比如同一订单的消息发到同一个 Partition,由 Key 决定分区。
分区机制:为什么分区,怎么分区
分区的意义
分区是 Kafka 并行扩展的底座:
- 写入并行:多个 Partition 可以分布在不同 Broker 上,多个 Producer 同时写入不同分区,I/O 压力分散
- 消费并行:同一个 Consumer Group 内,每个 Partition 只被一个 Consumer 消费,Consumer 数量上限 = Partition 数量
- 顺序保证:消息在 Partition 内有序,跨 Partition 无序
分区策略
Producer 决定消息发到哪个 Partition,常见策略:
| 策略 | 说明 | 适用场景 |
|---|---|---|
| 指定 Partition | 直接指定 partition 编号 | 特殊业务需求,如强制某个分区 |
| Key Hash | partition = hash(key) % numPartitions | 同一 Key 的消息发到同一分区,保证局部有序 |
| Round Robin | 轮询分配,Key 为空时默认使用 | 无需顺序保证,均匀分布 |
| 自定义分区器 | 实现 Partitioner 接口 | 业务定制逻辑,如按用户 ID 分区 |
常见坑:Key 设计不当导致分区倾斜。比如 Key 是”省份”,大部分用户集中在北上广,某个 Partition 堆积严重,其他 Partition 空闲。排查时用 kafka-consumer-groups --describe 查看各 Partition 的 LAG(积压量),如果集中在一两个 Partition,说明 Key 分布不均。
Partition 数量怎么定
不是越多越好。Partition 过多时:
- Broker 端:每个 Partition 对应多个文件(LogSegment),文件句柄耗尽,PageCache 污染
- 客户端:Consumer 需要维护更多连接和拉取请求
- ZooKeeper/KRaft:更多元数据,Controller 选举和 Rebalance 开销增加
经验值:单个 Broker 上 Partition 总数不超过 2000,单 Topic Partition 数量按吞吐量估算:
Partition 数量 ≈ 目标吞吐量 / 单 Consumer 最大吞吐量
如果单 Consumer 能处理 50MB/s,目标吞吐 500MB/s,Partition 数量至少 10。
ISR 机制:副本同步与可靠性保证
Kafka 的副本机制不追求”所有副本同步成功才返回”,而是通过 ISR(In-Sync Replica)动态维护与 Leader 保持同步的 Follower 集合。
ISR 的工作方式
每个 Partition 有多个 Replica,其中一个是 Leader,其他是 Follower:
- Leader:负责所有读写请求
- Follower:被动从 Leader 拉取数据,同步写入本地日志
- ISR:与 Leader 保持同步的 Follower 集合,只有 ISR 中的副本才能参与选举
同步判断标准:Follower 的 LEO(Log End Offset,日志末端位移)与 Leader 的 LEO 差距在 replica.lag.time.max.ms(默认 30 秒)以内。超时则踢出 ISR,恢复后重新加入。
acks 配置:生产者可靠性级别
Producer 发送消息时可以指定 acks 参数:
| acks 值 | 行为 | 可靠性 | 性能 |
|---|---|---|---|
| 0 | Producer 不等待任何 ACK | 最低,可能丢消息 | 最高 |
| 1 | Leader 写入成功后返回 ACK | 中等,Leader 故障可能丢消息 | 中等 |
| all/-1 | ISR 所有副本写入成功后返回 ACK | 最高,需配合 min.insync.replicas | 最低 |
acks=all 不是绝对可靠,需要配合:
min.insync.replicas >= 2:ISR 中至少 2 个副本写入成功- 副本数 >= 3:至少 3 个副本,1 Leader + 2 Follower
- ISR 不能只剩 1 个副本(说明其他副本都掉线了)
HW 和 LEO:消费者的可见边界
- LEO(Log End Offset):日志末端位移,下一条等待写入的消息位置
- HW(High Watermark):高水位,ISR 所有副本都已同步的消息位置
消费者只能拉取 HW 之前的消息,HW 之后的消息对消费者不可见。Leader 更新 HW 的时机:ISR 所有副本的 LEO 都达到某个位置后,HW 才推进到该位置。
这带来了一个问题:消息在 Leader 写入成功,但 HW 未推进时 Leader 故障,新 Leader 可能没有这条消息,消费者看不到,“已写入成功”的消息丢失。这是 Kafka 在可用性和一致性之间的取舍:允许短暂不一致,换取高吞吐。
消费者组与 Rebalance:消费模型的核心
消费者组的消费模型
Consumer Group 是 Kafka 实现消费扩展和负载均衡的核心机制:
- 同一个 Group 内:每个 Partition 只被一个 Consumer 消费,避免重复消费
- 不同 Group 之间:独立消费,实现”广播”效果(一个消息被多个 Group 各消费一次)
消费模型分为两种:
| 模型 | 说明 | 实现 |
|---|---|---|
| 队列模型 | 一个消息只被一个消费者处理 | 单 Consumer Group |
| 发布/订阅模型 | 一个消息被多个消费者各自处理 | 多 Consumer Group,各自独立消费 |
Offset 管理:消费进度存储
消费者的消费进度(Offset)有两种存储方式:
| 存储方式 | 说明 | 特点 |
|---|---|---|
| __consumer_offsets(Broker 端) | Kafka 内部 Topic,默认方式 | 支持自动提交、手动提交 |
| 本地存储 | Consumer 本地维护 Offset | 灵活控制,但故障时可能丢失 |
提交 Offset 的时机是关键:
- 自动提交:
enable.auto.commit=true,定期自动提交,可能提交未处理完的 Offset,消息丢失 - 手动提交:业务处理完成后手动提交
commitSync()或commitAsync()
推荐手动提交,确保业务处理成功后再提交 Offset。commitAsync() 性能好但可能失败,commitSync() 会阻塞等待成功,可根据场景选择。
Rebalance:消费者的重新分配
当 Consumer 加入、退出或崩溃时,触发 Rebalance,重新分配 Partition。Rebalance 是 Kafka 消费者组的”STW(Stop-the-world)“问题:
Rebalance 期间,所有 Consumer 停止消费,直到分配完成。
Rebalance 流程:
- Coordinator(Broker 上的 Group Coordinator)检测到 Consumer 失败或加入
- 发送 JoinGroup 请求,所有 Consumer 响应
- 选择 Group Leader,Leader 制定分配方案
- 发送 SyncGroup 请求,通知所有 Consumer 分配结果
- Consumer 开始按新分配消费
Rebalance 的常见触发原因:
- Consumer 宕机或主动退出
- Consumer 处理时间过长,超过
max.poll.interval.ms,被踢出 Group - Partition 数量变化(新增 Partition)
- Consumer Group 数量变化
Rebalance 优化:
- 合理设置
session.timeout.ms(默认 10s)和max.poll.interval.ms(默认 5min) - 使用 Cooperative Rebalance(增量再均衡):只重新分配变化的 Partition,不停止所有 Consumer(Kafka 2.4+)
- 避免消费者处理时间过长,单次
poll()处理的消息数控制在合理范围
Exactly-Once:从 At-Least-Once 到精确一次
Kafka 默认是 At-Least-Once:消息至少被消费一次,可能重复消费。业务需要自己做幂等。Kafka 0.11+ 引入了 Exactly-Once 语义的机制。
幂等生产者(Idempotent Producer)
单 Producer 单 Partition 内的幂等:Producer 分配一个 PID(Producer ID),每条消息带序列号,Broker 检测重复序列号并拒绝写入。
启用方式:
enable.idempotence=true
幂等生产者能解决单 Producer 单 Partition 的重复写入问题,但跨 Partition 不幂等,需要事务机制。
事务(Transactions)
跨 Partition 的原子写入:一组消息要么全部写入成功,要么全部失败。适用于”消费-处理-生产”链路,比如从 Topic A 消费,处理后写入 Topic B,整个过程原子化。
事务流程:
- Producer 初始化事务 ID(
transactional.id) - 开启事务
beginTransaction() - 发送消息到多个 Partition
- 发送 Offset 到 __consumer_offsets(如果是消费-处理-生产链路)
- 提交事务
commitTransaction()或回滚abortTransaction()
消费者需要设置 isolation.level=read_committed,只读取已提交的事务消息。
注意:Kafka 的事务主要针对跨 Partition 原子写入,不适合与外部数据库(MySQL)做分布式事务联动,且缺乏服务端事务回查机制。这是 RocketMQ 事务消息的强项。
常见误区
| 误区 | 真实情况 |
|---|---|
| Kafka 消息不丢 | acks=all + min.insync.replicas>=2 + 副本数>=3 + 手动提交 Offset,仍可能丢(Leader 故障时 HW 未推进) |
| Partition 越多越好 | Partition 过多导致文件句柄耗尽、PageCache 污染、Rebalance 开销增加 |
| Kafka 天生有序 | 只有 Partition 内有序,跨 Partition 无序 |
| Rebalance 无感 | Rebalance 会 STW,所有 Consumer 停止消费 |
| Kafka 事务 = 分布式事务 | Kafka 事务是跨 Partition 原子写入,不能与外部 DB 做分布式事务 |
项目判断
引入 Kafka 时,先回答三个问题:
- 数据量级:日均数据量是否达到 TB/PB ?是否需要喂给 Flink/Spark?
- 可靠性要求:是否允许偶尔丢消息?是否需要事务消息?
- Topic 数量:Topic 数量是否超过 50?是否需要延迟消息?
如果数据量大、允许偶尔丢消息、没有延迟消息需求 → Kafka 合适。如果涉及金融交易、需要事务回查、Topic 数量多 → RocketMQ 更合适。
线上排查积压时:
- 用
kafka-consumer-groups --describe查看 LAG,定位积压 Partition - 检查是否分区倾斜(Key 分布不均)
- 检查 Consumer 处理耗时,是否超过
max.poll.interval.ms - 紧急扩容 Consumer(数量不能超过 Partition 总数)
- 积压太大时,先转储消息,让 Consumer 追上进度,再离线补偿