消息队列01:Kafka

从分区机制、ISR 副本同步、消费者组 Rebalance 到 Exactly-Once 语义,拆解 Kafka 高吞吐背后的设计取舍和常见坑点。

字数 2780 阅读时长 ≈ 8 分钟 2026-7-13 2026-7-13
消息队列01:Kafka

Kafka 不是”消息队列”这个分类里最老的产品,但它把”高吞吐日志管道”这件事做到了极致。很多项目引入 Kafka,是因为日志量太大、埋点太多、或者要把数据喂给 Flink/Spark 做实时计算。

这篇文章不列所有参数,而是围绕四个核心机制:分区、ISR、消费者组、Exactly-Once,把 Kafka 的设计逻辑和常见坑点讲清楚。

为什么是 Kafka,而不是别的

先定位清楚:Kafka 的核心场景是海量数据日志/事件流的管道,追求极致吞吐量,适合 ETL、日志收集、事件溯源、流式计算。它不是为”金融交易消息”设计的,那是 RocketMQ 的主场。

两者的取舍很直接:

维度KafkaRocketMQ
存储模型分区式存储,每 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 Hashpartition = 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 值行为可靠性性能
0Producer 不等待任何 ACK最低,可能丢消息最高
1Leader 写入成功后返回 ACK中等,Leader 故障可能丢消息中等
all/-1ISR 所有副本写入成功后返回 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 流程:

  1. Coordinator(Broker 上的 Group Coordinator)检测到 Consumer 失败或加入
  2. 发送 JoinGroup 请求,所有 Consumer 响应
  3. 选择 Group Leader,Leader 制定分配方案
  4. 发送 SyncGroup 请求,通知所有 Consumer 分配结果
  5. 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,整个过程原子化。

事务流程:

  1. Producer 初始化事务 ID(transactional.id
  2. 开启事务 beginTransaction()
  3. 发送消息到多个 Partition
  4. 发送 Offset 到 __consumer_offsets(如果是消费-处理-生产链路)
  5. 提交事务 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 时,先回答三个问题:

  1. 数据量级:日均数据量是否达到 TB/PB ?是否需要喂给 Flink/Spark?
  2. 可靠性要求:是否允许偶尔丢消息?是否需要事务消息?
  3. Topic 数量:Topic 数量是否超过 50?是否需要延迟消息?

如果数据量大、允许偶尔丢消息、没有延迟消息需求 → Kafka 合适。如果涉及金融交易、需要事务回查、Topic 数量多 → RocketMQ 更合适。

线上排查积压时:

  • kafka-consumer-groups --describe 查看 LAG,定位积压 Partition
  • 检查是否分区倾斜(Key 分布不均)
  • 检查 Consumer 处理耗时,是否超过 max.poll.interval.ms
  • 紧急扩容 Consumer(数量不能超过 Partition 总数)
  • 积压太大时,先转储消息,让 Consumer 追上进度,再离线补偿