Kafka 架构深度
概述
Kafka 是分布式的高吞吐消息队列。架构核心:多个 Broker 组成集群,Topic 分分区,分区有副本。Controller 负责集群管理,ISR 保障副本一致性。本文讲透 Kafka 集群架构与核心机制。
一、Kafka 定位与架构
1.1 为什么用 Kafka
核心特性:
高吞吐(顺序写 + 页缓存)
高可用(多副本 + ISR)
持久化(落盘 + 可重放)
生态(流处理/数据集成)
应用:
日志收集、消息解耦
数据管道、流处理底座1.2 整体架构
Producer
↓ 写入
Kafka 集群(多个 Broker)
├── Controller(集群管理)
├── Broker 1(分区 Leader/Follower)
├── Broker 2
└── Broker 3
↓ 拉取
Consumer(消费组)| 角色 | 职责 |
|---|---|
| Producer | 生产消息 |
| Broker | 存储与转发 |
| Controller | 集群元数据管理 |
| Consumer | 消费消息 |
| ZooKeeper/KRaft | 元数据存储(新版本内置) |
二、Broker 与集群
2.1 Broker 职责
| 职责 | 说明 |
|---|---|
| 分区存储 | 管理本地分区副本 |
| 消息读写 | 处理请求 |
| 副本同步 | Follower 拉取数据 |
| 心跳上报 | 与 Controller 通信 |
Broker 内部模块:
请求处理(Acceptor/Processor)
日志存储(LogManager)
副本管理(ReplicaManager)
网络层(Selector)2.2 集群部署
典型部署:
3 台 Broker(最小高可用)
副本因子 2 或 3
Controller 自动选主
元数据存储演进:
ZooKeeper → KRaft(内置,去 ZK)| 版本 | 元数据 | 说明 |
|---|---|---|
| 2.x 前 | ZooKeeper | 依赖 ZK |
| 3.x+ | KRaft | 内置 Raft,简化运维 |
三、Controller 详解
3.1 Controller 职责
Controller 是集群的"大脑",负责:
| 职责 | 说明 |
|---|---|
| 分区 Leader 选举 | 分区副本故障时选新 Leader |
| 元数据管理 | 维护分区/副本状态 |
| 分区重分配 | 迁移副本 |
| Broker 上下线处理 | 感知节点变化 |
选举时机:
Broker 宕机 → 其 Leader 分区重新选举
分区副本变化 → 触发 Leader 选举
新副本上线 → 触发均衡3.2 Controller 选举
选举机制(KRaft):
Raft 协议选举 Controller
多数派同意 → 成为 Leader
旧 Leader 故障 → 重新选举
元数据日志复制:
所有变更写入日志
Follower 同步
保证元数据一致3.3 Controller 处理流程
处理链路:
监听 Broker 变化(ZNode/元数据日志)
→ 更新集群状态
→ 对受影响分区发起 Leader 选举
→ 广播新元数据给所有 Broker
→ Broker 更新本地缓存四、Topic 与分区
4.1 逻辑模型
Topic(主题)
└── Partition(分区,可并行)
├── Replica 0(Leader)
├── Replica 1(Follower)
└── Replica 2(Follower)| 概念 | 说明 |
|---|---|
| Topic | 逻辑队列 |
| Partition | 物理存储单元 |
| Replica | 分区副本 |
| Leader | 提供读写 |
| Follower | 同步备份 |
4.2 分区的作用
分区带来的能力:
并行:不同分区可由不同消费者消费
扩展:分区可分布到多 Broker
顺序:分区内有序
分区数确定:
吞吐需求 / 消费并行度
不能无脑多(文件句柄/选举开销)4.3 分区分配策略
分区分配原则:
均匀分布到各 Broker
副本不能在同一 Broker
优先考虑机架感知五、ISR 与 Leader 选举
5.1 ISR 概念
ISR(In-Sync Replicas):与 Leader 保持同步的副本集合。
同步定义:
副本追上 Leader 的 HW(高水位)
延迟在阈值内(replica.lag.time.max.ms)
只有 ISR 内的副本才有资格
被选举为 Leader5.2 LEO / HW / ISR 关系
| 概念 | 含义 |
|---|---|
| LEO | 日志末端偏移量(每个副本自己) |
| HW | 高水位(ISR 最小 LEO) |
| ISR | 同步副本集合 |
| AR | 全部分区副本集合 |
偏移量关系:
HW = ISR 中最小的 LEO
消费者只能读到 HW 之前的数据
HW 由 Leader 推进5.3 Leader 选举
正常选举:
从 ISR 中选(元数据中记录)
优先选副本序号最小的存活副本
ISR 全挂(极端):
允许 unclean 选举 → 选不在 ISR 的副本
(可能丢数据,换取可用性)| 配置 | 说明 |
|---|---|
| unclean.leader.election.enable | 是否允许非 ISR 副本选主 |
六、日志分段与存储
6.1 日志分段
分区日志 = 多个 Segment 文件
Segment 由 .log + .index + .timeindex 组成
Segment 达到大小/时间阈值滚动| 文件 | 作用 |
|---|---|
| .log | 消息数据 |
| .index | 稀疏偏移量索引 |
| .timeindex | 时间戳索引 |
6.2 偏移量管理
offset 语义:
Producer 写入 → 分配连续 offset
Consumer 消费 → 记录消费位置
分区内单调递增
消费者位移:
存 Kafka 内部 Topic(__consumer_offsets)
定期提交七、生产消费模型
7.1 生产流程
Producer → 分区器(选分区)→ 批次缓冲 → 网络发送 → Broker 落盘7.2 消费模型
Consumer Group:
组内消费者瓜分分区(1 分区 1 消费者)
组间互不影响(各自独立消费)
消费方式:
拉取模型(主动拉,控制速率)
位移提交(自动/手动)
Rebalance(分区重新分配)八、Kafka 高性能原理
8.1 为什么快
| 机制 | 说明 |
|---|---|
| 顺序写 | 磁盘顺序追加 |
| 页缓存 | 利用 OS 缓存 |
| 零拷贝 | sendfile 直接发 |
| 批处理 | 批量发送/拉取 |
| 分区并行 | 横向扩展 |
零拷贝原理:
传统:磁盘 → 内核 → 用户 → 内核 → 网卡(4 次拷贝)
Kafka:磁盘 → 内核 → 网卡(2 次拷贝,sendfile)8.2 网络线程模型
线程模型:
Acceptor 线程:接受连接
Processor 线程:解析请求
KafkaRequestHandler 线程池:处理请求
队列解耦:Processor → 请求队列 → Handler九、常见问题
9.1 分区数怎么定
估算:
目标吞吐 = 分区数 × 单分区吞吐
消费并行度 ≤ 分区数
建议:按峰值吞吐 + 预留余量9.2 Controller 故障影响
Controller 短时间故障:
分区 Leader 选举暂停
已有读写不受影响
元数据变更延迟
应对:
快速恢复(多副本选主)
监控 Controller 状态9.3 副本同步落后
Follower 落后原因:
网络慢/机器性能差
分区数过多
配置阈值过小
处理:
监控 ISR 收缩
落后副本自动踢出 ISR
排查瓶颈后重新同步十、小结
Kafka 架构核心是分区复制模型 + Controller 集群管理。理解 ISR/LEO/HW 才能掌握数据可靠性,理解日志分段与零拷贝才能理解性能。把"高吞吐、高可用、可持久化"背后的机制串起来,是掌握 Kafka 的关键。