编程知识 编程知识CODING KNOWLEDGE
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Kafka

Kafka

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)元数据获取机制

  • 生产者不会被动接收集群推送,元数据是惰性刷新

  • 刷新时机:

    1. 定时刷新:默认 5 分钟强制拉取最新元数据。

    2. 发送消息刷新:本地缓存元数据失效/分区变更时,发送消息瞬间触发刷新。

  • 分区扩容后,生产者不会立刻感知,短时间内只会使用旧分区列表。

(2)默认分区路由策略 DefaultPartitioner

  • 携带 Keyhash(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 副本同步完成,可靠性最高,性能最低,金融核心业务使用。

九、面试高频易错点总结(必背)

  1. Kafka 是拉模式;消费者主动拉取消息。

  2. 分区是并发单元,副本是高可用单元,副本不提升吞吐量。

  3. 分区内有序、跨分区无序;想要全局有序只能单分区。

  4. Rebalance 仅消费者存在,生产者只有元数据惰性刷新。

  5. 同组消费者不能超分区数,多余消费者空闲。

  6. 自动提交丢消息、手动提交重复消费;生产一律手动提交+幂等。

  7. Offset 只记录位点,不删除消息,支持消息重放。

  8. Follower 不处理读写,仅同步数据。

  9. 分区扩容会打乱 Key 有序性,生产禁止随意扩容。

  10. 不同消费组可消费同一分区,实现广播;同组绝对互斥。

思考

offset存储与 __consumer_offsets 内部主题

  1. 消费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高可用。
  1. offset提交可靠性逻辑
    消费者提交offset的内部生产者默认 acks=all(-1):
    只有Leader写入,并且全部ISR副本同步完成,才返回commit成功给消费者。
  • 如果已经收到commit成功:即使随后__consumer_offsets的Leader宕机,ISR选出的新Follower副本拥有完整offset记录,不会丢失offset,不会因此重复消费。
  • 如果offset提交过程中,还未返回成功,__consumer_offsets Leader宕机:本次提交视作失败,新Leader没有这条offset记录,重启/rebalance后读取旧offset,引发重复消费。
  1. 重要区分
    业务topic的partition发生Leader切换(业务分区Leader宕机),不会影响offset存储,不会造成offset丢失与重复消费。只有__consumer_offsets自身集群故障才会影响offset。

__consumer_offsets 深入理解

  1. 消费者提交offset的底层行为
    消费者客户端内部维护一个隐藏的Producer实例;提交offset本质就是使用该内部Producer,将offset元数据作为消息发送到__consumer_offsets

消费者客户端同时具备Consumer与Producer双重身份;内部producer默认 acks=all(-1),保障offset提交可靠性。

  1. __consumer_offsets 本质
    本质就是Kafka内部普通Topic,具备partition、replica、Leader/Follower、ISR;offset以消息形式持久化在该topic日志中。
    消费组长期不活跃,offset记录会被日志清理策略删除,再次消费会触发auto‑offset‑reset逻辑。

  2. 专属Broker配置(仅topic首次自动创建生效)

  • offsets.topic.num.partitions:内部topic分区数,默认50;创建完成后修改参数不会自动扩容。
  • offsets.topic.replication.factor:内部topic副本因子,生产建议3;单机测试环境需改为1。

以上参数与业务topic分区、副本配置相互独立,互不影响。业务topic分区与副本在创建topic时指定。

返回列表