面试专题

Kafka 是如何实现高并发、高性能、高可用的

从分区并行、顺序写、Page Cache、批处理、零拷贝、副本与 ISR 等机制,系统理解 Kafka 高吞吐与高可用架构。

难度:深入更新:2026-07-24

Kafka 是如何实现高并发、高性能、高可用的

面试回答

Kafka 的高并发主要依靠分区实现横向扩展:一个 Topic 被拆成多个 Partition,分区 Leader 分散在不同 Broker 上,生产者可以并行写入,消费者组也可以按分区并行消费。高性能主要来自追加写和顺序 I/O、操作系统 Page Cache、生产与消费两端的批处理、批量压缩,以及通过 sendfile 实现的零拷贝传输。高可用则依靠多副本、Leader/Follower、ISR、故障选主和 KRaft Controller Quorum;生产端再配合 acks=allmin.insync.replicas 和幂等生产,消费端通过 Consumer Group、Offset 与重新分配实现故障恢复。Kafka 的核心设计是在分区粒度上并行,在单分区内部维持有序日志,再用副本机制保障故障后的连续服务。

一句话总结:

分区提供并行度,顺序日志、批处理和零拷贝提供吞吐量,副本、ISR 与故障选主提供可用性。

详细讲解

Kafka 的三个能力不是彼此孤立的:

  • 分区既是并发执行和横向扩展的单位,也是副本复制与故障恢复的单位。
  • 批处理既减少网络请求,也把大量小写入转化为更高效的顺序写。
  • Page Cache 既提升读写性能,也能配合副本机制避免为了可靠性而对每条消息执行同步刷盘。
  • 副本增强了可靠性,但同时会增加网络、磁盘与确认延迟,因此需要在吞吐、延迟和可靠性之间取舍。

Kafka 高并发、高性能和高可用实现机制总览

一、高并发:通过分区横向扩展

1. Partition 是 Kafka 的并行单元

一个 Topic 可以拆成多个 Partition:

Topic: order-events
├── Partition 0 → Leader 在 Broker A
├── Partition 1 → Leader 在 Broker B
└── Partition 2 → Leader 在 Broker C

不同分区可以由不同 Broker 同时处理,因此整体吞吐不再受限于单台机器。增加 Broker 并合理增加分区后,存储容量、网络带宽和读写能力都可以横向扩展。

Kafka 只保证单分区内有序,不保证一个 Topic 的所有分区之间全局有序。这是它获得并行能力的重要前提。

2. 生产者直接找到分区 Leader

生产者会缓存集群元数据,根据消息 Key、显式分区或默认分区策略选择 Partition,然后直接向该分区的 Leader 发送数据,中间不需要额外的统一路由节点。

常见分区方式:

  • 指定 Key:相同 Key 通常进入同一分区,适合保证同一订单、用户或账户内的消息顺序。
  • 不指定 Key:默认分区策略会在适当时机选择分区,并尽量形成较大的批次。
  • 自定义分区器:按租户、地区或业务规则路由。

如果 Key 分布不均,大量消息集中到少数分区,就会出现热点分区。此时即使 Broker 很多,整体吞吐仍会被热点 Leader 限制。

3. Consumer Group 按分区并行消费

同一个 Consumer Group 中,一个 Partition 在同一时刻只分配给一个 Consumer;一个 Consumer 可以负责多个 Partition。

例如:

4 个 Partition + 2 个 Consumer
→ 每个 Consumer 负责约 2 个 Partition

4 个 Partition + 4 个 Consumer
→ 每个 Consumer 负责约 1 个 Partition

4 个 Partition + 6 个 Consumer
→ 最多只有 4 个 Consumer 能获得 Partition

因此,一个消费组的有效并行度上限通常受 Partition 数量约束。简单增加 Consumer,而不增加 Partition,不一定能提高吞吐量。

4. Broker 内部通过线程分工处理请求

Broker 不会用一个线程串行完成所有网络和磁盘工作。网络线程负责接收请求、解析协议和返回响应,请求处理线程负责执行 Produce、Fetch 等具体逻辑,后台线程负责日志清理、副本同步等任务。

这种职责分离让网络连接管理、业务请求处理和后台维护能够并行执行,避免某类工作完全阻塞其他工作。

二、高性能:把随机小操作变成顺序批量操作

1. 顺序追加写

Kafka 的 Partition 本质上是一条只追加的日志。新消息写到日志末尾,不需要像通用数据库那样频繁进行随机位置更新。

旧数据 ────────────────→ 日志末尾追加新批次
offset 0  1  2  3  4  5  6 ...

顺序 I/O 能充分利用磁盘和操作系统的预读、合并写能力。日志达到一定大小或时间后会切分为新的 Segment,旧 Segment 可以按保留策略清理或压缩,不影响当前日志继续追加。

2. 分段日志与稀疏索引

每个 Partition 由多个日志 Segment 组成,常见文件包括:

00000000000000000000.log
00000000000000000000.index
00000000000000000000.timeindex
  • .log 保存消息批次。
  • .index 保存相对 Offset 到日志物理位置的稀疏映射。
  • .timeindex 保存时间戳到相对 Offset 的稀疏映射。

查找消息时,Kafka 先定位 Segment,再通过稀疏索引找到接近目标的位置,最后在日志中继续顺序查找。索引不需要记录每一条消息,因此能控制索引体积和内存占用。

3. 充分利用 Page Cache

Kafka 不把全部消息封装成大量 Java 堆对象长期缓存,而是依赖操作系统 Page Cache:

Producer
→ Broker 写日志文件
→ 数据先进入 Page Cache
→ 操作系统异步回写磁盘

这样做有几个好处:

  • 利用操作系统成熟的预读和回写机制。
  • 减少 JVM 堆内存占用和 GC 压力。
  • 热数据可以直接从 Page Cache 读取。
  • Broker 重启后,操作系统缓存仍可能保留部分热数据,不必完全重新构建 JVM 内缓存。

Kafka 的可靠性主要依赖多副本,而不是对每条消息都执行一次同步 fsync。如果强制每条消息同步刷盘,吞吐和延迟都会明显恶化。

4. 端到端批处理

批处理贯穿生产、Broker 存储和消费三段链路:

多条消息
→ Producer 聚合成 Record Batch
→ 一次网络请求发送
→ Broker 将 Record Batch 连续追加到分区日志
→ Consumer 一次 Fetch 多条消息

它可以摊薄:

  • 网络往返成本。
  • 系统调用成本。
  • 协议头和请求头开销。
  • 磁盘写入次数。
  • 消息压缩成本。

这里的“Broker 追加批次”不是把多条消息合并成一条业务消息,也不承诺底层一定只发生一次系统调用。它表示 Kafka 的生产者和存储格式都以 Record Batch 为基本单位:

Record Batch
├── 批次头:baseOffset、长度、时间戳、压缩类型、CRC 等
├── Record 1
├── Record 2
└── Record 3

Producer 会先按 Partition 在内存中积累消息,形成 Record Batch。Broker 收到后,对批次进行校验,为其中的消息确定 Offset,然后把整段连续数据追加到对应 Partition 的当前日志 Segment 末尾。相比每条消息都单独发请求、单独写日志,批次追加可以用更少的请求和更大的连续 I/O 摊薄固定开销。

一个 ProduceRequest 还可以携带多个 Partition 的 Record Batch,但 Broker 最终仍然分别追加到各个 Partition 的日志中。因此:

批次是网络传输和日志写入的基本单位,Partition 才是日志组织、顺序保证和副本复制的基本单位。

生产者常见相关参数:

batch.size=16384
linger.ms=5
compression.type=lz4

batch.sizelinger.ms 体现的是吞吐与延迟之间的取舍:适度等待可以积累更大的批次,但等待过长会增加低流量场景的发送延迟。参数值应根据消息大小、流量和延迟目标压测确定,不能机械照搬示例。

5. 批量压缩

Kafka 对一个 Record Batch 进行整体压缩,而不是分别压缩每条消息。相似消息中的字段名和公共内容可以获得更好的压缩率。

压缩后的批次会以压缩形式写入日志并传输给消费者,从而降低:

  • Broker 磁盘占用。
  • 生产者到 Broker 的网络流量。
  • 副本复制流量。
  • Broker 到消费者的网络流量。

压缩会消耗 CPU,因此需要根据资源瓶颈选择 lz4zstdsnappygzip 等算法。网络或磁盘是瓶颈时,压缩往往能换取更高吞吐。

6. 零拷贝

消费者拉取日志数据时,普通路径可能需要:

磁盘 → Page Cache → 用户空间 → Socket Buffer → 网卡

Kafka 在适用条件下利用 Linux sendfile,让数据从 Page Cache 直接进入网络发送路径,减少用户态和内核态之间的复制与上下文切换:

Page Cache → Socket / 网卡

这就是常说的零拷贝。它主要优化 Broker 向消费者传输已有日志数据的过程。

Kafka sendfile 零拷贝与传统文件传输路径对比

传统路径需要先把数据从内核 Page Cache 复制到 Broker 的用户空间,再从用户空间写回内核 Socket Buffer;sendfile 让内核直接组织 Page Cache 到网络 Socket 的传输,Broker 不再把消息内容完整搬入 JVM 用户空间。

所以“零拷贝”更准确的含义是:

减少或绕过 Page Cache 与用户空间之间不必要的数据复制,而不是数据在磁盘、内存和网卡之间一次都不移动。

需要注意,Kafka 官方文档明确说明,启用 SSL 时数据要经过用户态 TLS 处理,当前不会使用这条 sendfile 路径。因此不能笼统地说所有 Kafka 网络传输都一定是零拷贝。

7. Consumer Pull 与长轮询

Kafka 由消费者主动 Pull 数据。消费者可以根据自己的处理能力决定拉取速度,处理不过来时,未消费消息会暂时保留在 Kafka 中,并表现为 Consumer Lag 增长,而不是由 Broker 持续 Push 直至压垮消费者。

Consumer Lag 表示消费者的消费进度落后于 Partition 最新位置的程度。监控消费组时,通常可以近似理解为:

Consumer Lag
= Partition Log End Offset
- Consumer Group 已提交 Offset

例如:

Partition 下一条待写位置:10000
消费组下一条待消费位置:9400

Lag = 10000 - 9400 = 600

这表示当前大约还有 600 个 Offset 的数据尚未被该消费组处理。Lag 不是消息丢失,而是消费积压:

生产速度 > 消费速度
→ Lag 持续增大

消费速度 > 生产速度
→ Consumer 逐步追赶
→ Lag 逐步减小

Lag 不能只看某一时刻的绝对值,还要观察增长趋势、积压持续时间和业务允许的最大延迟。同样数量的 Lag,在每秒几条和每秒几十万条的 Topic 中代表的时间延迟完全不同。

如果 Lag 长期增长,可能是消费者处理变慢、下游故障、分区热点、Consumer 数量不足或频繁重平衡。若积压时间超过 Topic 的数据保留期限,旧消息可能已经被清理,消费者就不能再依靠 Kafka 把这部分数据完整追回。

Fetch 请求支持批量拉取和长轮询:

fetch.min.bytes=1
fetch.max.wait.ms=500

Broker 暂时没有足够数据时,可以等待一段时间再响应,避免消费者无数据时频繁空轮询;流量较高时,又可以一次返回较大的批次。

三、高可用:副本、ISR 与故障转移

1. 每个 Partition 可以有多个副本

假设副本因子为 3:

Partition 0
├── Leader:Broker A
├── Follower:Broker B
└── Follower:Broker C

生产和普通消费请求由 Leader 处理,Follower 分别从 Leader 拉取并复制日志。多个 Follower 的同步是彼此独立的,不存在先同步 B、再由 B 同步 C 的固定顺序。

副本应尽量分散在不同 Broker;具备机架感知配置时,还可以跨机架放置,降低单机或单机架故障带来的影响。

2. ISR 维护可安全选主的副本集合

ISR 是 In-Sync Replicas,即当前与 Leader 保持同步的副本集合。Follower 故障或长时间落后时会被移出 ISR,重新追上后可以再次加入。

正常:ISR = [A, B, C]
C 落后:ISR = [A, B]

Leader 故障后,Controller 会从符合条件的副本中选择新的 Leader。以 ISR 为基础进行选主,可以避免把明显缺少最新日志的副本直接作为新的数据源。

3. acksmin.insync.replicas

常见的高可靠组合是:

replication.factor = 3
min.insync.replicas = 2
acks = all

含义是生产者等待当前 ISR 的确认,并且 ISR 至少要有两个副本,否则 Broker 拒绝写入。

这个配置体现了 CAP 取舍:副本不足时拒绝写入,会降低部分可用性,但可以避免在只剩单副本时继续写入并承担更高的数据丢失风险。

4. KRaft Controller Quorum 保证控制面可用

KRaft Controller Quorum 是什么

KRaft 是 Kafka 自己实现的元数据管理机制,用来替代早期依赖的 ZooKeeper。集群中会有若干个承担 controller 角色的节点,这些 Controller 组成一个基于 Raft 共识协议的 Controller Quorum,即控制器仲裁组。

它们维护的是集群元数据,而不是普通 Topic 的业务消息:

  • Broker 注册、存活状态和节点信息。
  • Topic 与 Partition 的定义。
  • 每个 Partition 的副本分配、Leader 和 ISR。
  • 配置、配额以及其他集群级元数据。

这些变化会写入 KRaft 的集群元数据日志,并复制给 Quorum 中的其他 Controller。一次元数据变更获得多数 Controller 确认后才算提交,因此单个 Controller 故障不会导致已提交元数据丢失。

常见部署是 3 个或 5 个 Controller:

3 个 Controller → 多数派为 2,可容忍 1 个故障
5 个 Controller → 多数派为 3,可容忍 2 个故障

如果存活 Controller 失去多数派,集群不能安全提交新的元数据变更。此时已有数据读写是否还能短暂继续,取决于具体操作和现有元数据状态,但创建 Topic、Partition Leader 调度等控制面操作会受到影响。

Active Controller 是什么

Controller Quorum 在同一时刻会选出一个 Leader,Kafka 的文档和运行语境中通常称它为 Active Controller。其他 Controller 作为 Follower 复制元数据日志,随时准备接管。

Controller A:Active Controller
├── 接收和处理元数据变更
├── 管理 Broker 注册与故障
├── 触发 Partition Leader 选举
└── 把变更写入元数据日志

Controller B、C:Follower Controllers
├── 复制元数据日志
└── Active Controller 故障后参与重新选举

Active Controller 不是所有业务消息的入口,也不负责替 Broker 保存普通 Topic 数据。生产者和消费者仍然直接与 Partition Leader 所在的 Broker 通信。

需要明确区分两个层面:

层面主要角色保存什么负责什么
数据面Broker、Partition Leader/Follower普通 Topic 的业务消息生产、消费、副本复制
控制面KRaft Controller Quorum集群元数据日志节点管理、元数据变更、Leader 调度

例如 Broker A 突然故障:

Active Controller 感知 Broker A 失联
→ 找出受影响的 Partition
→ 从符合条件的副本中选择新 Leader
→ 将选举结果写入元数据日志
→ 多数 Controller 确认提交
→ 向相关 Broker 发布新元数据
→ 客户端刷新元数据并访问新 Leader

如果 Active Controller 自身故障,剩余 Controller 会通过 Raft 重新选出新的 Active Controller。新节点已经复制了已提交的元数据日志,因此可以继续管理集群。

5. Consumer Group 自动接管故障消费者的分区

消费者通过心跳维持组成员身份。如果某个 Consumer 崩溃或超时,Consumer Group 会重新分配它原来负责的 Partition,让其他 Consumer 接管。

消费进度保存在已提交的 Offset 中,因此新 Consumer 可以从相应位置继续处理:

Consumer A 故障
→ 触发重新分配
→ Partition 交给 Consumer B
→ B 从已提交 Offset 继续消费

重平衡期间可能产生短暂停顿,所以高可用不等于无感切换。合理设置会话、心跳、处理时长,并避免频繁扩缩容,有助于减少重平衡影响。

四、三种能力如何协同

目标核心机制解决的问题主要代价
高并发多 Partition、多 Broker、Consumer Group并行生产、存储和消费分区与元数据管理成本
高性能顺序写、Page Cache、批处理、压缩、零拷贝降低 I/O、复制和网络开销批处理延迟与 CPU 消耗
高可用多副本、ISR、Controller Quorum、故障转移节点故障后继续提供服务副本存储、网络与确认延迟

可以把整个链路理解为:

消息按 Partition 分流
→ Producer 批量发送
→ Leader 顺序追加到日志
→ Follower 并行复制
→ 满足确认条件后返回
→ Consumer Group 按 Partition 批量拉取

五、为什么不能只靠增加分区提升性能

分区是扩展能力的基础,但分区不是越多越好。分区数量增加会同时带来:

  • 更多日志目录、Segment 和文件句柄。
  • 更多副本复制任务。
  • 更大的集群元数据。
  • 更高的 Leader 选举和故障恢复成本。
  • Consumer Group 重平衡开销。
  • 单 Key 顺序范围更难调整。

合理做法是根据目标吞吐量、单分区压测能力、消费者并行度、副本数和未来增长空间估算,而不是一次创建大量分区。

一个简化估算思路:

分区数 ≥ max(
  目标生产吞吐 / 单分区生产吞吐,
  目标消费吞吐 / 单消费者单分区吞吐,
  期望消费并行度
)

最终仍需要在接近生产环境的机器、网络、副本配置和消息大小下进行压测。

六、常见误区

误区一:Kafka 使用磁盘,所以一定比内存消息队列慢

Kafka 采用顺序追加、Page Cache 和批处理,大量读写实际在操作系统缓存中完成。性能取决于访问模式,而不是简单由“使用磁盘”决定。

误区二:零拷贝让数据完全不发生复制

零拷贝是减少不必要的用户态复制,并非物理意义上的一次复制都没有;而且启用 SSL 时不会走 Kafka 文档描述的 sendfile 路径。

误区三:消费者越多,消费速度一定越快

同一消费组内,一个 Partition 同时只能交给一个 Consumer。Consumer 数超过 Partition 数后,多出的 Consumer 不会获得分区。

误区四:副本越多,高可用越高且没有代价

更多副本会增加存储、网络复制和写确认成本。副本数应根据故障域和可靠性目标确定,常见三副本不代表所有场景都必须相同。

误区五:acks=all 就代表所有配置副本都已写入

acks=all 等待的是当前 ISR,而不是所有配置副本。要避免 ISR 只剩一个副本时仍然成功写入,需要配合 min.insync.replicas

七、排查性能和可用性问题时看什么

生产端

  • 发送吞吐、请求延迟、错误率和重试次数。
  • 批次大小、压缩率和缓冲区等待。
  • 消息 Key 是否造成分区流量倾斜。

Broker

  • 各 Broker、Partition 和 Leader 的流量是否均衡。
  • 磁盘利用率、磁盘延迟、网络带宽和请求队列。
  • ISR 收缩、未充分复制分区和离线分区。
  • Page Cache 命中效果和系统是否发生 Swap。

消费端

  • Consumer Lag 及其增长速度。
  • 单条消息处理耗时和批次处理耗时。
  • Consumer 数量与 Partition 数量是否匹配。
  • 重平衡次数、心跳超时和 Offset 提交失败。

参考资料

DISCUSSION

评论与讨论

留下你的想法