Kafka
开篇:Kafka 为什么能扛住每秒百万级消息?
如果你在一个大型快递公司工作,每天要处理上百万件包裹。普通的方式是一件一件登记、分拣、投递,效率可想而知。但如果你把包裹按区域分成几十条流水线,每条线并行运转,效率就完全不同了。
Kafka 就是这样一个"超级快递中心"。它是 LinkedIn 开源、后来进入 Apache 基金会的分布式流处理平台,以每秒百万级消息的吞吐量闻名于世。无论是日志采集、实时数仓还是事件驱动架构,Kafka 都是当之无愧的主力军。
接下来,我们从核心概念出发,逐步拆解 Kafka 为什么这么快、这么可靠。
一、核心概念:邮局模型
理解 Kafka 最好的方式,是把它想象成一个巨型邮局。
Topic——信箱
你要给朋友寄信,首先要写上收件地址。在 Kafka 里这就是 Topic(主题)。Topic 是消息的逻辑分类,比如"订单消息""日志消息"各对应一个 Topic。生产者往指定的 Topic 里塞消息,消费者从 Topic 里取消息。
Partition——分拣通道
一个信箱如果只有一个口,高峰期就会排长队。Kafka 的解法是把一个 Topic 拆成多个 Partition(分区),每个 Partition 就是一条独立的流水线。消息在每个 Partition 内部是严格有序的,但不同 Partition 之间没有顺序保证。
Topic: order-events
├── Partition-0 → msg0, msg1, msg2, ...
├── Partition-1 → msg0, msg1, msg2, ...
└── Partition-2 → msg0, msg1, msg2, ...为什么有了 Topic 还要 Partition?软件领域有句老话:任何问题都可以加一个中间层来解决。在 Topic 的基础上再细粒度地划分一层,主要带来三个好处:
- 提升吞吐量——多条流水线并行处理,每个 Partition 可以由不同消费者独立消费
- 负载均衡——消息均匀分散到多个 Broker 节点上
- 水平扩展——增加 Partition 数量就能线性提升集群的并发处理能力和存储容量
Offset——信件编号
每条进入 Partition 的消息都有一个唯一的递增编号,叫 Offset(偏移量)。消费者通过 Offset 知道自己"读到哪了",下次从这个位置继续。
一个关键细节:Consumer 记录的 Offset 指向的是下一条要消费的消息,而不是最后一条已消费的。比如已经处理了 0~5 号消息,那 Offset 应该是 6,不是 5。
消费者需要把自己的消费进度上报给 Kafka(称为"提交 Offset"),避免消费者崩溃重启后不知道从哪继续。提交方式有自动和手动两种,后面详细介绍。
Consumer Group——邮递员小队
Kafka 用 Consumer Group(消费者组) 来实现消息的负载均衡。同一个 Group 内的多个 Consumer 会"瓜分"Topic 的 Partition,每个 Partition 只会被组内一个 Consumer 消费。这保证了消息不会被重复处理(在同一个 Group 内)。
如果 Consumer 数量少于 Partition 数量,部分 Consumer 会消费多个 Partition。反过来,如果 Consumer 多于 Partition,多出来的 Consumer 会闲置,因为默认情况下一个 Partition 只能被一个 Consumer 消费。
Kafka 4.0(2025 年 3 月发布)引入了共享组(Shared Group),允许多个消费者消费同一个 Partition 并支持逐条确认。这打破了传统限制,让扩容变得更灵活。
整体架构
Kafka 集群由多个 Broker 组成。每个 Broker 上可以存储多个 Partition 的副本(Replica),其中一个是 Leader(负责读写),其余是 Follower(负责同步数据)。Leader 挂了,Follower 可以自动接替,实现高可用。
之前集群的协调工作由 ZooKeeper 完成(Broker 注册、Leader 选举、元数据管理等)。从 Kafka 2.8 开始引入 KRaft 协议,到 Kafka 4.0 已默认使用 KRaft,不再依赖 ZooKeeper。好处是运维更简单、故障恢复更快、扩展性更好。
二、生产者流程:消息是怎么发出去的?
当你调用 producer.send(msg) 时,消息并不是直接飞到 Broker,中间还经过了好几道工序。
2.1 发送链路
整个过程涉及两个线程:
- Main 线程:负责初始化 Producer、执行发送逻辑。消息经过拦截器、序列化器、分区器后,被放入 RecordAccumulator 缓冲区
- Sender 线程:后台运行,定期从缓冲区取出消息批次,通过 NetworkClient 发送给对应 Partition 的 Leader Broker
每一步的职责:
- 拦截器——在发送前后对消息做定制处理(加日志、修改消息体等)
- 序列化器——把 Key 和 Value 对象转成字节数组,便于网络传输
- 分区器——决定消息发往哪个 Partition。指定了 Key 就按 Key 的 hash 值路由;没有 Key 则轮询(Round-robin)
- RecordAccumulator——消息不会一条一条发,而是先攒在缓冲区里,满足
batch.size(批次大小)或linger.ms(等待时长)条件后打包成批次,大大减少网络往返次数 - Sender 线程——负责真正的网络发送,处理确认和重试
2.2 同步发送 vs 异步发送
// 同步发送——阻塞等待结果
SendResult result = producer.send(record).get();
// 异步发送——回调方式(推荐)
producer.send(record, (metadata, exception) -> {
if (exception != null) {
log.error("发送失败,topic={}", record.topic(), exception);
// 重试逻辑
}
});同步发送简单直接,但吞吐量低。异步发送配合回调,既能保证高吞吐,又能感知失败进行重试。
注意:直接调用 producer.send(msg) 是异步的,方法立即返回并不代表消息发送成功。如果不做任何处理,发送失败你根本不知道。所以务必使用带 callback 的方式。
2.3 ACK 机制
Broker 收到消息后要不要确认?确认到什么程度?这由 acks 参数控制:
| acks 值 | 含义 | 可靠性 | 吞吐量 |
|---|---|---|---|
| 0 | 发完就走,不等确认 | 最低 | 最高 |
| 1 | Leader 写入即确认 | 中等 | 中等 |
| -1 (all) | 所有 ISR 副本写入后确认 | 最高 | 最低 |
生产环境中如果要保证消息不丢失,建议配置:
acks=-1
retries=3
retry.backoff.ms=3002.4 怎么发到同一个 Partition?
顺序消费场景需要保证相关消息落到同一个 Partition。三种方式:
方式一:指定 Partition
// 直接指定 partition=0
ProducerRecord<String, String> record =
new ProducerRecord<>("topic", 0, null, "message");方式二:指定 Key(推荐)
相同 Key 的消息会被 hash 到同一个 Partition:
ProducerRecord<String, String> record =
new ProducerRecord<>("topic", null, orderId, "message");方式三:自定义 Partitioner
实现 Partitioner 接口,在 partition() 方法中编写自定义路由逻辑,然后在 Producer 配置中指定:
props.put("partitioner.class",
"com.example.CustomPartitioner");三、消费者与消费者组
3.1 消费模型
Kafka 的消费是 Pull(拉取) 模式——消费者主动去 Broker 拉消息,而不是 Broker 推过来。好处是消费者可以按自己的速度来,不会被大流量冲垮。Kafka 还支持批量拉取(一次 poll 多条消息),进一步提升效率。
3.2 Offset 提交策略
消费者需要告诉 Kafka 自己消费到哪了。有两种方式:
自动提交
enable.auto.commit=true
auto.commit.interval.ms=5000每隔 5 秒自动提交当前 poll 到的最大 Offset。优点是简单,缺点是消息可能还没处理完就提交了——如果此时消费者崩溃,这些消息就"丢"了(Kafka 以为你已经消费了)。
手动提交
enable.auto.commit=false处理完消息后显式调用提交方法。commitSync() 同步等待提交成功,commitAsync() 异步提交不阻塞。手动提交更安全,能确保"至少一次"语义。
| 对比维度 | 自动提交 | 手动提交 |
|---|---|---|
| 实现复杂度 | 简单 | 需要自己管理 |
| 可靠性 | 较低,可能丢消息或重复 | 高,处理完才提交 |
| 适用场景 | 日志、统计等容错场景 | 金融、交易等关键场景 |
| 性能 | 较高,提交频率低 | 稍低,提交更频繁 |
3.3 批量消费不丢消息的正确姿势
批量消费时最容易踩坑。两个常见错误:
错误一:用自动提交——消息还没处理完,Offset 已经自动提交了。
错误二:在 finally 块里手动提交——不管成功失败都提交了,等于自动提交。
正确做法是全部成功才提交:
@KafkaListener(topics = "my-topic",
containerFactory = "batchListenerFactory")
public void listen(List<ConsumerRecord<?, ?>> records,
Acknowledgment ack) {
CompletionService<Boolean> cs =
new ExecutorCompletionService<>(executor);
List<Future<Boolean>> futures = new ArrayList<>();
// 1. 提交所有任务到线程池
for (ConsumerRecord<?, ?> record : records) {
futures.add(cs.submit(() -> {
processMessage(record); // 业务处理
return true;
}));
}
// 2. 检查每个任务的结果
boolean allSuccess = true;
for (int i = 0; i < records.size(); i++) {
try {
if (!cs.take().get()) {
allSuccess = false;
break;
}
} catch (Exception e) {
allSuccess = false;
break;
}
}
// 3. 全部成功才提交 Offset
if (allSuccess) {
ack.acknowledge();
}
// 失败则不提交,消息会重投(需做好幂等)
}这么做消息不会丢了,但会带来消息重投——只要有一个失败,整批消息都会重投。不过这比丢消息好得多,消费端做好幂等就行了。
3.4 重平衡(Rebalance)
当消费者组的成员发生变化(新加入、退出、崩溃),或者订阅的 Topic/Partition 数量变化时,Kafka 会触发重平衡——重新把 Partition 分配给消费者。
重平衡期间,所有消费者会暂停消费(类似 JVM GC 的 STW),这是 Kafka 使用中最令人头疼的问题之一。
重平衡可能带来的问题:
- 消费暂停——吞吐量短暂归零
- 重复消费——Offset 还没提交就被重新分配给其他消费者
- 消息堆积——暂停期间消息在 Broker 积压
优化手段:
- CooperativeStickyAssignor(渐进式重平衡)——Kafka 2.4 引入。不再让所有消费者释放全部 Partition,而是只调整需要变更的部分,其他消费者继续消费不受影响
partition.assignment.strategy=\
org.apache.kafka.clients.consumer.CooperativeStickyAssignor静态成员——设置
group.instance.id,消费者短暂离线再回来时保持原有分配,不触发重平衡Kafka 4.0 下一代协议——将分区分配逻辑从客户端移到服务端,消费者可以独立重平衡,不必暂停整个组
渐进式重平衡的工作原理:
重平衡前:
ConsumerA -> P0, P1, P2
ConsumerB -> P3, P4, P5
新增 ConsumerC:
第一阶段(部分撤销):
ConsumerA 释放 P2(P0、P1 继续消费)
ConsumerB 释放 P5(P3、P4 继续消费)
第二阶段(重新分配):
ConsumerA -> P0, P1
ConsumerB -> P3, P4
ConsumerC -> P2, P5相比传统方式"全员停工重排队",渐进式只影响需要调整的 Partition。
四、副本机制:数据安全的保障
4.1 ISR 机制
Kafka 每个 Partition 都有多个副本(Replica),分为一个 Leader 和若干 Follower。所有的读写请求都走 Leader,Follower 只负责从 Leader 同步数据。
ISR(In-Sync Replicas) 是与 Leader 保持同步的副本集合。判断标准是 replica.lag.time.max.ms(默认 10 秒):如果某个 Follower 的同步落后 Leader 超过这个时间,就会被踢出 ISR。
只有 ISR 中所有副本都确认收到消息后,这条消息才算"已提交",消费者才能看到。
早期版本(0.9.x 之前)用的是
replica.lag.max.messages(基于消息条数),但在瞬间高并发时容易误判——比如 Leader 瞬间收到几万条消息,所有 Follower 来不及同步就会被全部踢出 ISR。所以后来改成了基于时间的判断,只要在限定时间内追上来就行。
4.2 高水位(HW)与 LEO
这两个概念用来管理消费者的可见范围和数据一致性:
Partition 中的消息:
|--- 已提交(消费者可见) ---|--- 未提交(仅 Leader 有) ---|
0 1 2 3 4 5 6 7 8
↑ HW ↑ LEO- LEO(Log End Offset)——日志末尾偏移量,指向下一条待写入消息的位置
- HW(High Watermark)——高水位,标识所有 ISR 副本都已同步到的位置。消费者只能拉取 HW 之前的消息
HW 的两个作用:
- 让消费者知道哪些消息已经安全提交可以消费了
- 在副本切换时作为数据一致性的参考点
4.3 Leader Epoch
Leader 切换时,新旧 Leader 的数据可能不一致。Kafka 引入 Leader Epoch(递增整数)来标识每任 Leader 的"任期"。
每次副本切换时,新 Leader 会增加自己的 Epoch。新 Leader 检查旧 Leader 的 Epoch 和 HW,只有在旧 Leader 的 Epoch <= 新 Leader 的 Epoch,且旧 Leader 的 HW <= 新 Leader 的 HW 时,才会接受旧 Leader 的数据。
这个机制避免了数据回滚问题,确保新 Leader 不会接受不属于当前任期的消息。
4.4 Leader 选举
Kafka 中有两类选举:
Partition Leader 选举:当某个 Partition 的 Leader 挂了,Kafka 会从该 Partition 的 ISR 中选一个 Follower 提升为新 Leader。
Controller 选举:整个 Kafka 集群有一个 Controller 节点,负责管理所有 Partition 的副本分配和 Leader 选举。Controller 本身也需要选举——多个 Broker 竞争创建 ZooKeeper 上的 /controller 临时节点,先到先得。KRaft 模式下改用 Raft 协议选举。
保证数据安全的关键配置:
replication.factor=3
min.insync.replicas=2
unclean.leader.election.enable=falseunclean.leader.election.enable=false 表示不允许非 ISR 成员当选 Leader。虽然可能短暂不可用,但能保证数据不丢。
五、高性能设计:Kafka 为什么这么快?
Kafka 的高性能是全链路优化的结果,我们从发送、存储、消费三个阶段来看。
5.1 发送端优化
| 技术 | 原理 |
|---|---|
| 批量发送 | 多条消息打包成 batch,减少网络往返次数 |
| 异步发送 | 不等 ACK 就继续发下一批,主线程不阻塞 |
| 消息压缩 | 支持 GZIP/Snappy/LZ4/ZStd,减少网络传输量 |
| 并行发送 | 不同 Partition 对应不同 Broker,天然并行 |
5.2 存储端优化
| 技术 | 原理 |
|---|---|
| 顺序写磁盘 | 消息以追加方式写入日志文件,避免随机寻道。顺序写的速度接近内存 |
| Page Cache | 利用操作系统页缓存,热数据留在内存中,读写极快 |
| 零拷贝 | 使用 sendfile() 系统调用,数据从磁盘直达网卡,不经过用户态拷贝 |
| 稀疏索引 | 不为每条消息建索引,而是每隔 N 条建一个索引点。兼顾查找速度和索引体积 |
| 分段存储 | 日志文件按大小切割成 Segment,方便过期数据清理,也便于查找 |
5.3 消费端优化
| 技术 | 原理 |
|---|---|
| 消费者组 | 多消费者并行消费不同 Partition |
| 批量拉取 | 一次 poll 拉取多条消息,减少网络开销 |
| 零拷贝 | 读取时同样受益于 sendfile() |
5.4 数据存储结构详解
Kafka 的存储是一个从逻辑到物理的层级映射:Topic -> Partition -> Segment。
/kafka-logs/
├── my-topic-0/ # Partition 0 对应一个文件夹
│ ├── 00000000000000000000.log # 消息数据(Segment 1)
│ ├── 00000000000000000000.index # Offset 索引
│ ├── 00000000000000000000.timeindex # 时间索引
│ ├── 00000000000000005000.log # Segment 2(文件名=首条消息 Offset)
│ ├── 00000000000000005000.index
│ ├── 00000000000000005000.timeindex
│ └── leader-epoch-checkpoint
├── my-topic-1/ # Partition 1
│ └── ...
└── ...每个 Segment 包含三类文件:
.log——消息数据,顺序追加写入.index——Offset 索引,存储相对 Offset 与物理位置的稀疏映射.timeindex——时间索引,根据时间戳快速定位消息
写入流程:生产者消息到达后,追加到当前活跃 Segment 的 .log 文件末尾。只有当消息被写入磁盘(或 Page Cache),写入才算成功。
读取流程:消费者指定 Offset -> 根据 Segment 文件名(文件名就是首条消息的 Offset)快速定位到对应文件 -> 用 .index 文件做二分查找定位物理位置 -> 从 .log 中顺序扫描读取消息 -> 通过零拷贝发送给消费者。
六、Exactly-Once 语义
Kafka 提供三种消息传递语义:
| 语义 | 含义 | 典型场景 |
|---|---|---|
| At most once | 最多一次,可能丢 | 不重要的日志 |
| At least once | 至少一次,可能重复(默认) | 大部分业务 |
| Exactly once | 精确一次,不丢不重 | 金融交易 |
Exactly-Once 依赖 Kafka 的事务机制和幂等生产者。
需要理解的关键点:Kafka 的事务消息保证的是一组消息的原子性——要么全部成功提交,要么全部回滚。这和 RocketMQ 不同,RocketMQ 的事务消息保证的是"本地事务 + 消息发送"的原子性。
事务消息实现
背后的关键组件:
- Transaction Coordinator(事务协调器)——管理事务状态,分配 Transaction ID
- Producer ID + Epoch——唯一标识生产者及其版本,防止因重启导致的重复消息
- Transaction Log(事务日志)——内部特殊 Topic,记录事务的开始、提交、回滚
整体流程遵循两阶段提交:
生产者代码示例:
props.put("transactional.id", "my-tx-id");
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.initTransactions();
try {
producer.beginTransaction();
producer.send(new ProducerRecord<>("topic1", "k1", "v1"));
producer.send(new ProducerRecord<>("topic2", "k2", "v2"));
producer.commitTransaction();
} catch (KafkaException e) {
producer.abortTransaction(); // 异常时回滚
}消费端需要设置隔离级别,只读取已提交的事务消息:
isolation.level=read_committed七、常见面试题精选
1. Kafka 为什么不能 100% 保证消息不丢失?
Kafka 默认提供 At least once 语义,但即使做了各种配置,极端场景下仍然无法 100% 保证。从三个环节分析:
- 生产者:异步发送可能在回调前崩溃;重试也可能持续失败;Producer 进程突然被 kill
- Broker:消息写入 Page Cache 后、刷盘前宕机;Leader 同步到 Follower 前宕机;即使设
log.flush.interval.messages=1也有微小窗口 - 消费者:自动提交 Offset 后、处理完成前崩溃
acks=-1 + min.insync.replicas>1 可以大幅降低丢失概率,但极端场景(如多副本同时宕机、机房断电)仍有风险。要做到真正的"不丢",需要引入本地消息表或分布式事务作为兜底。
2. 消息丢了,最可能的原因是什么?
| 环节 | 常见原因 | 频率 |
|---|---|---|
| 生产者 | acks=0 或 1 | 常见 |
| 生产者 | 未处理发送失败回调 | 常见 |
| 生产者 | 发送过程中 Producer 崩溃 | 偶发 |
| Broker | 异步刷盘时宕机 | 不常见 |
| Broker | Leader 未同步给 Follower 就挂了 | 不常见 |
| Broker | 消息超过保留时间(默认 72h)未消费 | 偶发 |
| 消费者 | 自动提交后消息未处理完 | 常见 |
| 消费者 | 批量消费时 finally 中直接 ack | 常见 |
防止丢失的最佳实践配置:
# 生产者
acks=-1
retries=3
retry.backoff.ms=300
# Broker
replication.factor=3
min.insync.replicas=2
unclean.leader.election.enable=false
# 消费者
enable.auto.commit=false3. 如何实现顺序消费?
Kafka 只保证 Partition 内有序。要实现业务层面的顺序消费:
- 只创建一个 Partition(简单粗暴,丧失并行能力,不推荐)
- 用相同的 Key 发送到同一个 Partition(推荐,比如用订单 ID 做 Key)
- 自定义 Partitioner 精确控制路由
4. 如何保证只消费一次?
需要多管齐下:
- 消费者组——同组内一条消息只会被一个消费者处理
- 手动提交 Offset——处理完才提交,避免跳过
- 客户端幂等——消费者端做好"一锁、二判、三更新"
- Exactly-Once 语义——使用事务保证消息消费和 Offset 提交的原子性
5. 什么时候该选 Kafka?
| 维度 | Kafka | RocketMQ | RabbitMQ |
|---|---|---|---|
| 吞吐量 | 百万级/秒 | 十万级/秒 | 万级/秒 |
| 延迟 | 毫秒级 | 毫秒级 | 微秒级 |
| 事务消息 | 消息原子性 | 本地事务+消息原子性 | 基本事务 |
| 延迟队列 | 不原生支持 | 原生支持 | 插件支持 |
| 语言 | Scala/Java | Java | Erlang |
| 最佳场景 | 日志、大数据、流处理 | 电商、金融 | 中小规模、复杂路由 |
6. Kafka 4.0 有什么重要变化?
Kafka 4.0 于 2025 年 3 月发布,几个关键变化:
- 默认 KRaft——不再依赖 ZooKeeper,简化运维
- 共享组——多个消费者可以消费同一个 Partition,支持逐条确认,突破了消费者数量限制
- 下一代重平衡协议——分区分配从客户端移到服务端,消费者可独立重平衡,大幅减少 STW
小结
Kafka 的核心设计可以用一句话概括:用 Partition 做并行,用副本做容错,用顺序写 + 零拷贝 + Page Cache 做极致性能。
理解了这些,再去看具体的配置参数和 API,就不会迷失在细节里。记住那个邮局模型:Topic 是信箱,Partition 是分拣通道,Offset 是信件编号,Consumer Group 是邮递员小队。一切围绕"如何又快又稳地分拣海量包裹"展开。
最后,附一张关键配置速查表:
| 参数 | 推荐值 | 作用 |
|---|---|---|
acks | -1 | 所有 ISR 确认后才返回成功 |
retries | 3 | 发送失败自动重试次数 |
replication.factor | 3 | 每个 Partition 的副本数 |
min.insync.replicas | 2 | 最小同步副本数 |
unclean.leader.election.enable | false | 禁止非 ISR 成员当选 Leader |
enable.auto.commit | false | 关闭自动提交,使用手动提交 |
isolation.level | read_committed | 只读取已提交的事务消息 |
使用消息队列的三大核心价值:解耦(模块间通过消息通信,互不依赖)、异步(耗时操作放队列后台处理,主流程立即返回)、削峰填谷(高峰请求先进队列排队,下游按能力慢慢消费)。选择 Kafka 还是其他 MQ,取决于你最看重哪个维度——吞吐量选 Kafka,事务可靠性选 RocketMQ,灵活路由选 RabbitMQ。
最后一个提醒:Kafka 的消息不是消费完就删除的(这和 RabbitMQ 不同),而是基于
log.retention.hours(默认 72 小时)或log.retention.bytes的保留策略定期清理。如果消费者在保留期内没有消费完消息,过期的消息会被删除——这也是一种容易被忽视的"消息丢失"场景。