Kafka 核心原理、机制、风险点、面试考点完整版
一、Kafka 核心整体架构
1. 核心组件
-
Broker:Kafka 服务节点,集群多个 Broker 组成整体服务,负责存储消息、处理读写请求、持久化日志数据。
-
KRaft 控制器:新版本替代 Zookeeper,负责管理集群元数据、Broker 状态、分区 Leader 选举、分区分配。
-
Topic 主题:消息逻辑分类载体,用于业务区分消息,本身不存储数据,数据分散在各个分区。
-
Partition 分区:Topic 物理拆分单元、Kafka 唯一并发单元。分区内消息有序,跨分区无序。
-
Replica 副本:分区数据备份统称,分为 Leader、Follower,用于高可用、故障转移。
-
Leader 副本:分区主副本,所有生产者写入、消费者读取请求只走 Leader。
-
Follower 副本:分区从副本,不处理客户端请求,只主动同步 Leader 数据,故障时可升级为 Leader。
-
ISR 同步副本集合:与 Leader 数据保持同步的副本列表,只有 ISR 内副本可竞选新 Leader。
-
Offset 偏移量:分区内消息唯一序号;消费 Offset 记录消费位点,保存在内部主题
__consumer_offsets。 -
Producer 生产者:推送消息,主动路由分区,写入分区 Leader。
-
Consumer 消费者:拉模式主动拉取消息,处理业务后提交 Offset。
-
ConsumerGroup 消费组:逻辑消费者集群,通过 group-id 区分,实现负载均衡与消息广播。
2. 整体运行机制
-
Kafka 是 消费者拉模式,消费者主动向 Broker 拉取消息。
-
消息写入磁盘持久化,消费后不会删除消息,依靠过期时间/文件大小自动清理。
-
Offset 仅标记消费位点,可重置,支持历史消息重放。
二、Topic 与 Partition 核心机制
1. 分区核心特性
-
一个 Topic 可以包含多个 Partition 分区。
-
分区内消息有序,跨分区无序。
-
分区数量决定吞吐量与消费并发上限,副本数只决定高可用,不提升并发。
2. 生产者如何获取分区信息 & 路由规则
(1)元数据获取机制
-
生产者不会被动接收集群推送,元数据是惰性刷新。
-
刷新时机:
-
定时刷新:默认 5 分钟强制拉取最新元数据。
-
发送消息刷新:本地缓存元数据失效/分区变更时,发送消息瞬间触发刷新。
-
-
分区扩容后,生产者不会立刻感知,短时间内只会使用旧分区列表。
(2)默认分区路由策略 DefaultPartitioner
-
携带 Key:
hash(key) % 分区数,相同 Key 固定同一分区,保证 Key 维度消息有序。 -
Key 为 Null:轮询分发,消息均匀打散所有分区,吞吐量最高,无序。
-
支持手动指定分区、自定义分区器。
(3)分区扩容风险
-
已有 Topic 扩容分区后,hash 取模结果改变,相同 Key 会进入不同分区,顺序彻底打乱。
-
生产环境禁止随意扩容已有业务 Topic 分区。
三、副本机制(Leader / Follower / ISR)
1. 核心规则
-
一个 Partition 同一时刻 有且仅有一个 Leader,其余全部是 Follower。
-
角色是 Partition 维度,不是 Broker 维度:一台 Broker 可以同时是多个分区的 Leader、多个分区的 Follower。
2. Leader 副本作用
-
处理所有生产者写入、消费者读取请求。
-
负责消息落地磁盘、维护 ISR 同步副本集合。
3. Follower 副本作用
-
不处理任何客户端读写请求。
-
主动拉取 Leader 日志数据,同步备份。
-
Leader 宕机时,ISR 内 Follower 可升级为新 Leader,实现故障转移。
4. ISR 同步副本集合
-
定义:与 Leader 数据同步、无明显滞后的副本集合。
-
只有 ISR 内副本具备 Leader 竞选资格。
-
acks=all代表:等待所有 ISR 副本同步完成,才返回写入成功,可靠性最高。
四、消费组 ConsumerGroup 核心机制
1. 消费组定义
-
相同 group-id 的消费者属于同一个消费组,不同 group-id 完全隔离。
-
消费 Offset 归属:group-id + topic + partition,不属于消费者实例。
2. 同组消费规则(重中之重)
-
同一个消费组内,一个 Partition 同一时刻只能被一个消费者消费。
-
消费者实例数量 不能超过分区数,多余消费者空闲,无法提升并发。
-
组内多消费者实现 负载均衡,一条消息只会被组内一个消费者消费。
3. 不同消费组规则
-
不同消费组可同时消费同一个分区,各自维护独立 Offset。
-
用于实现 消息广播:一份消息多业务消费。
4. auto-offset-reset 机制
仅消费组无历史 Offset 时生效:
-
earliest:从头消费所有历史消息。
-
latest:只消费后续新消息。
五、Rebalance 重平衡机制(纯消费者机制)
1. 核心结论
- 生产者没有 Rebalance,Rebalance 只属于消费者消费组。
2. Rebalance 触发条件
-
消费组成员变更:消费者启动、宕机、下线、心跳超时被踢出组。
-
订阅 Topic 分区数量变更(分区扩容)。
-
消费者订阅 Topic 列表变更。
3. Rebalance 影响
-
重平衡期间消费暂停,频繁 Rebalance 严重影响业务吞吐量。
-
是 Kafka 重复消费的核心诱因之一。
4. 生产者对应行为(易混点)
- 分区扩容、Leader 切换时,生产者只会 惰性刷新元数据,无 Rebalance。
六、Offset 提交机制 & 核心风险
1. 核心区别:生产者 ACK / 消费者 Offset
-
生产者 acks:Broker 确认消息持久化、副本同步结果,保障发送可靠性。
-
消费者无 acks 机制,依靠 Offset 提交 上报消费位点,不删除消息。
2. 自动提交(enable-auto-commit=true)
-
机制:后台线程定时提交(默认5秒),与业务处理进度无关。
-
致命风险:消息丢失
- Offset 已提交,批量消息部分未处理,消费者宕机,未处理消息永久丢失。
-
适用场景:日志、监控埋点,允许少量消息丢失的低可靠业务。
3. 手动提交(enable-auto-commit=false)
-
机制:业务代码处理完成后手动提交 Offset。
-
风险:消息重复消费
- 业务处理完成,未提交 Offset 时发生宕机 / Rebalance,重启后重复消费。
-
优势:绝对不会丢失消息。
-
生产必选:订单、支付、交易等核心业务。
4. 最终选型结论
-
高可靠业务:关闭自动提交,手动提交 + 业务幂等。
-
低可靠、日志类业务:可使用自动提交。
-
无论哪种模式,都可能出现重复消费,幂等是业务必备。
七、重复消费专项考点(面试高频)
1. 同消费组会不会多个消费者同时消费一个分区?
不会。Kafka 分配机制保证同一时刻一个分区仅属于组内一个消费者,不存在并发争抢重复消费。
2. 那重复消费从哪来?
根源:业务处理成功,Offset 未提交,触发 Rebalance / 消费者宕机
-
分区被重新分配给组内其他消费者。
-
新消费者从旧 Offset 开始拉取,导致已处理消息重复消费。
3. 解决手段
-
业务实现幂等(唯一键去重、状态机控制)。
-
优化消费超时时间,减少无故 Rebalance。
-
业务处理完成尽快提交 Offset,避免堆积大量未提交位点。
八、生产者 acks 三档机制与可靠性
-
acks=0:不等待 Broker 应答,吞吐量最高,极易丢消息。
-
acks=1(默认):Leader 写入成功即返回,兼顾性能与可靠性,Leader 宕机可能丢消息。
-
acks=all(-1):等待所有 ISR 副本同步完成,可靠性最高,性能最低,金融核心业务使用。
九、面试高频易错点总结(必背)
-
Kafka 是拉模式;消费者主动拉取消息。
-
分区是并发单元,副本是高可用单元,副本不提升吞吐量。
-
分区内有序、跨分区无序;想要全局有序只能单分区。
-
Rebalance 仅消费者存在,生产者只有元数据惰性刷新。
-
同组消费者不能超分区数,多余消费者空闲。
-
自动提交丢消息、手动提交重复消费;生产一律手动提交+幂等。
-
Offset 只记录位点,不删除消息,支持消息重放。
-
Follower 不处理读写,仅同步数据。
-
分区扩容会打乱 Key 有序性,生产禁止随意扩容。
-
不同消费组可消费同一分区,实现广播;同组绝对互斥。
思考
offset存储与 __consumer_offsets 内部主题
- 消费offset位点(
group‑id + topic + partition)持久化存储在Kafka内部topic:__consumer_offsets,不存储在业务topic的分区中。
- __consumer_offsets 也是普通topic,拥有自己的分区、副本、Leader/Follower、ISR集合。
- 消费者提交offset本质:内部生产者向
__consumer_offsets发送offset记录消息。 - 参数
offsets.topic.replication.factor控制该内部topic副本因子,生产建议≥3,保障offset高可用。
- offset提交可靠性逻辑
消费者提交offset的内部生产者默认 acks=all(-1):
只有Leader写入,并且全部ISR副本同步完成,才返回commit成功给消费者。
- 如果已经收到commit成功:即使随后
__consumer_offsets的Leader宕机,ISR选出的新Follower副本拥有完整offset记录,不会丢失offset,不会因此重复消费。 - 如果offset提交过程中,还未返回成功,
__consumer_offsetsLeader宕机:本次提交视作失败,新Leader没有这条offset记录,重启/rebalance后读取旧offset,引发重复消费。
- 重要区分
业务topic的partition发生Leader切换(业务分区Leader宕机),不会影响offset存储,不会造成offset丢失与重复消费。只有__consumer_offsets自身集群故障才会影响offset。
__consumer_offsets 深入理解
- 消费者提交offset的底层行为
消费者客户端内部维护一个隐藏的Producer实例;提交offset本质就是使用该内部Producer,将offset元数据作为消息发送到__consumer_offsets。
消费者客户端同时具备Consumer与Producer双重身份;内部producer默认 acks=all(-1),保障offset提交可靠性。
-
__consumer_offsets 本质
本质就是Kafka内部普通Topic,具备partition、replica、Leader/Follower、ISR;offset以消息形式持久化在该topic日志中。
消费组长期不活跃,offset记录会被日志清理策略删除,再次消费会触发auto‑offset‑reset逻辑。 -
专属Broker配置(仅topic首次自动创建生效)
- offsets.topic.num.partitions:内部topic分区数,默认50;创建完成后修改参数不会自动扩容。
- offsets.topic.replication.factor:内部topic副本因子,生产建议3;单机测试环境需改为1。
以上参数与业务topic分区、副本配置相互独立,互不影响。业务topic分区与副本在创建topic时指定。