Theory · Kafka · Internals

Kafka 底层原理与高吞吐

分布式提交日志(commit log)的性能四重奏:顺序写 + 页缓存 · 批量与端到端压缩 · 零拷贝 · 分区并行

一句话本质

Topic 分区 = 只增不改的日志段文件序列,读写全部 O(1),与总数据量无关

六大支柱

顺序追加 · page cache · sendfile 零拷贝 · 批量+压缩 · 分区并行 · 稀疏索引二分定位

一手验证

基于 Kafka 4.3 官方文档与 trunk 源码(2026-08),版本差异逐处标注

这份 deck 回答两个高频面试主题:Kafka 的底层存储/网络/复制原理,以及它高吞吐的技术细节。组织顺序:全景 → 设计哲学 → 存储与读写路径 → 零拷贝 → 网络 → 生产者 → 消息格式 → 消费者 → 复制 → 并行度 → 调优 → QA。所有数字与源码结论都对照 4.3 官方文档和 trunk 代码核对过。

Why Kafka Exists

先看麻烦:一次支付成功,8 个系统都要这条数据

具体场景:用户支付成功,风控、积分、库存、推送、报表、审计、搜索索引、对账——8 个下游系统都需要这条事件,而且要准实时。先看三个"看起来能行"的朴素方案为什么都不行。

方案 A:支付服务同步调用 8 个下游

一次支付要等 8 次 RPC:任一下游慢,主链路就慢;任一下游挂,支付就失败(或者你要自己写重试、降级、幂等)。更糟的是耦合——加第 9 个下游要改主服务的代码并重新发版

方案 B:写一张 MySQL 事件表,下游各自扫表

解耦了,但存储侧立刻崩:每个下游都扫表 + 更新 status 字段,是随机 IO 加行锁竞争;消费完还得 DELETE/归档(删了就再也重放不了);数据一多扫表就慢,吞吐上限很低。

方案 C:用传统 MQ(队列语义,消费即删)

解耦也做到了,但"取走即消失":8 个下游各要一份,就得复制 8 条队列;broker 还要为每条消息维护"已发送/已确认"状态,吞吐受限;出了问题想重放历史——数据已经没了。

四个绕不开的诉求:解耦(生产者不必知道谁在消费);② 多订阅(一份数据被 N 个下游按各自节奏各读一遍);③ 可重放(数据不能被"消费"掉);④ 吞吐(每秒百万级事件)。

Kafka 的答案:把它做成一个"只追加、可多人同时读、读了不删"的日志

只追加(append-only):生产者只往日志末尾写,不修改、不删除已有记录 —— 顺序写让磁盘跑满带宽。
每人一个游标:每个读者自己记住"我读到第几条"(一个整数 offset),互不干扰 —— 多订阅天然满足,谁快谁慢自己决定,一个人的进度不影响别人。
按时间保留而非消费即删:数据默认留 7 天,读多少次都在 —— 可重放天然满足,出事可以倒回去重算。
性能靠结构换:顺序写 + 页缓存 + 零拷贝,让读写几乎全是线性的,吞吐逼近磁盘与网络的物理极限。

换掉了什么,付出了什么

换来的是"一份数据、多个视角、随时重放";付出的是语义变简单了:没有传统的"消息确认/重投/优先级"那些花活,进度只是一个数字。这个取舍正是它快的根本原因——broker 不需要为每条消息记账

本 deck 的路线

先立全景与设计哲学(第 4–5 页)→ 给出"高吞吐六大支柱"总纲(第 6 页)→ 再按数据流向逐层拆:存储(第 7–8 页)· 零拷贝(第 9 页)· 网络(第 10 页)· 生产者(第 11–12 页)· 消费者(第 13 页)· 复制(第 14 页)· 分区(第 15 页)→ 最后速查与 QA。读完你应该能回答"Kafka 为什么这么快"以及"它的快是怎么一步步换来的"。

动机页:用"一次支付成功、8 个下游都要"的具体场景,让三个朴素方案(同步 RPC=耦合且脆弱、MySQL 事件表=随机 IO 且不可重放、传统 MQ=消费即删且要复制多份)依次失败,逼出四个真诉求(解耦/多订阅/可重放/吞吐),再落到 Kafka 的四个设计选择。关键论点:进度只是一个整数 offset,所以 broker 不必为每条消息记账——这是快的根本原因。禁止一上来就抛 topic/partition/ISR 三个名词。

Prerequisites & Glossary

先把词认全:下面每一页都会用到它们

术语一句话理解(先记住这个,细节后面展开)
消息
Message / Record
一条业务数据 = key + value + 时间戳 + 若干 header;本 deck 里"消息""记录"同义
主题
Topic
一类消息的命名流(像一张表、一个文件名),生产与消费都按名字寻址
分区
Partition
一个 topic 被切成若干条有序、只追加的子日志;保序与并行的最小单位
位移
Offset
消息在分区内的递增序号(0,1,2…);每个读者的消费进度就是"下一个要读的 offset"
生产者 / 消费者
Producer / Consumer
往日志末尾追加的客户端 / 按自己游标读的客户端;两端都是普通进程里的库
代理
Broker
一台 Kafka 服务器进程;若干台组成一个集群,每台存一部分分区
副本 / Leader / Follower
Replica
分区的多份拷贝;Leader 负责读写,Follower 主动拉取同步,故障时从追平的副本里选新 Leader
批次
Batch
若干条消息打包成的一块数据 —— 网络传输、磁盘追加、压缩、复制的共同单位
日志段
Segment
分区目录下的一个文件(默认满 1GB 就换新文件),文件名是该段首条消息的 offset
页缓存
Page Cache
操作系统拿空闲内存缓存磁盘数据的一块区域;读写都先经过它,进程重启它依然温热

如果你还不熟"顺序 IO / 零拷贝 / 页缓存"

先读这两篇再回来,本 deck 把这三件事当成已知:
OS · 零拷贝 → sendfile / SG-DMA 的内核机制
OS · 文件系统 → 文件、页缓存与顺序/随机 IO 的差距

一个最小心智模型(后面所有机制都是它)

把 Kafka 想成:一个只能往后追加、可以很多人同时读、读了不会消失的文件。
· 写 = 追加到文件末尾(顺序 IO)
· 读 = 从自己记住的位置往后读(一个整数游标)
· 删 = 按时间整段回收(不是按"消费过"删)
后面所有的性能设计,本质都在做两件事:让读写尽量保持顺序、让数据尽量少被拷贝几次

提前建立的三个预期

① 你会反复看到"与数据量无关":因为一切定位都是 O(1) 的二分 + 顺序扫描。
② 你会反复看到"批量":批越大,网络/磁盘/压缩的单位开销摊得越薄。
③ 你会反复看到"把活交给操作系统":缓存交给页缓存、拷贝交给 sendfile、刷盘时机交给内核。

阅读提示:术语不用背,遇到忘了的回来查这一页就行;真正需要背的是第 17–18 页那两张速查表(为什么快 / 必背数字 / 易错点)。
前置页:十个术语先定义再使用(message/topic/partition/offset/producer+consumer/broker/replica/batch/segment/page cache)。右侧给两条 OS 前置 deck 链接,并把全部机制收拢成"一个只能追加、可多人读、读了不删的文件"这一个意象。三个预期(与数据量无关 / 批量 / 交给操作系统)为后面六大支柱铺路。

Big Picture

一张图看懂 Kafka:分布式、分区的提交日志

上一页的答案:一个只追加、可多人读、读了不删的日志。这张图把它放大到集群规模——日志被切成多个分区散布多机,每份存多副本。

Kafka 架构全景图 生产者把批量压缩后的消息直接发给分区 Leader 所在的 Broker;每个分区有一个 Leader 和若干 Follower,Follower 主动拉取复制并构成 ISR;消费者组按分区拉取数据;KRaft quorum 负责集群元数据与选主,4.0 起取代 ZooKeeper。 PRODUCE FETCH 元数据·选主 BROKER CLUSTER · TOPIC T(3 分区 · 副本数 2+) CLIENT Producer ×N RecordAccumulator 批量·压缩·acks Broker 1 P0 · Leader P1 · Follower P2 · Follower 顺序日志文件 Broker 2 P0 · Follower P1 · Leader P2 · Follower 顺序日志文件 Broker 3 P0 · Follower P1 · Follower P2 · Leader 顺序日志文件 KRaft Quorum —— 集群元数据(Raft 日志)· 4.0 起移除 ZooKeeper CLIENT Consumer Group 分区↔消费者独占 offset 即消费状态 LEGEND Produce(直发 Leader) Fetch(零拷贝) Follower 拉取复制(ISR) Leader
先立骨架:Kafka 不是队列,是分布式提交日志。三个要点——① 生产者不经过任何路由层,直接发到分区 Leader;② 每个分区一个 Leader + N 个 Follower,Follower 用和消费者一样的方式拉数据,所以复制天然成批;③ 元数据由 KRaft(Raft 日志)管理,4.0 起彻底移除 ZooKeeper。消费组保证分区内一个消费者独占一个分区,消费状态只是一个 offset 整数。

Design Rationale

Don't fear the filesystem —— 用磁盘而不是堆

600 MB/s6×7200rpm SATA RAID-5 顺序写(官方文档)
100 kB/s同盘阵随机写 —— 差距 6000×
28–30 GB32GB 机器上可用的 page cache(全部空闲内存)

为什么不用堆内缓存

  • Java 对象开销高:数据入堆体积翻倍甚至更糟,实际存了 两份(进程缓存 + OS pagecache)
  • 堆越大 GC 越慢越碎;页缓存在堆外,零 GC 压力
  • 进程重启缓存冷;page cache 重启依然温热
  • 缓存一致性交给内核(read-ahead 预读 + write-behind 合并写)

为什么 O(1) 优于 BTree

  • BTree 操作 O(log N):一次操作可能多次磁盘寻道,且磁盘同一时刻只能一个寻道,并行受限
  • 缓存与磁盘混合访问下,BTree 性能随数据量超线性衰减(数据翻倍,慢超两倍)
  • 追加写 + 顺序读全 O(1):性能与数据量完全解耦 → 可以放心保留 7 天数据而不是消费即删

官方结论:"All data is immediately written to a persistent log on the filesystem without necessarily flushing to disk. In effect this just means that it is transferred into the kernel's pagecache." —— 写盘即写页缓存,fsync 默认不做。

这是 Kafka 一切性能设计的出发点:磁盘没有想象中慢——顺序 600MB/s 对随机 100kB/s,6000 倍。所以它反着设计:所有数据立刻进文件系统(其实是进页缓存),不维护堆内缓存。四个理由:JVM 对象开销、GC、重启冷缓存、内核一致性维护更高效。第二层是数据结构选择:不用 BTree 用纯追加日志,读写全 O(1),性能与数据量解耦,这才敢默认保留 7 天历史数据。引用句建议背下来。

Throughput Ingredients

高吞吐六大支柱(后面逐一展开)

1 顺序追加写

分区 = 只增不改的 segment 文件;随机小 IO 被批量合流成顺序大 IO,顺序写快随机写约 6000×

2 OS 页缓存

读写全走 page cache,堆外零 GC;消费者追上进度时磁盘零读,全部命中缓存

3 sendfile 零拷贝

消费路径 FileRecords.writeTo → transferTo,数据不过用户态;4 次拷贝降到 2 次

4 批量 + 端到端压缩

消息以 RecordBatch 为单位收发存储;producer 压缩 → broker 原样落盘 → consumer 解压

5 分区并行

分区是并行与顺序写的交换单元:多分区散布多 broker 多磁盘,吞吐近线性扩展

6 稀疏索引 + 二分

每 4KB 一条 .index 索引(mmap),二分定位 segment 与物理位置,读任何 offset 都是常数级

官方原话:"This simple optimization produces orders of magnitude speed up... turn a bursty stream of random message writes into linear writes."(批量把随机写洪峰变成线性写)

六大支柱是回答"Kafka 为什么快"的总纲:1、2 解决写盘慢;3 解决读网络传输;4 把小 IO 合成大 IO 并压网络;5 提供水平扩展;6 保证随机 offset 读取也是常数时间。面试时先报这六点,再按面试官追问展开到对应页。

Storage Model

Topic → Partition → Segment → Index

Kafka 存储模型与 offset 定位 分区目录内按 1GB 滚动生成多个 segment,文件名为该段首条 offset;.index 稀疏索引每 4KB 记一条 offset 到物理位置的映射并用 mmap 映射,读取时先二分定位 segment,再二分索引得到近似位置,最后顺序扫描。 my-topic-0/ 00000000000000000000.log 00000000000000000000.index 00000000000000000000.timeindex 00000000000000172453.log 00000000000000172453.index 000000000000003xxxxx.log ← 活跃段 segment 文件名 = 该段首条 offset (20 位补零,天然有序可二分) 删除按 segment 整段进行,不碰正在写的段 目标 offset = 5000 二分① 00000...00000.index 稀疏索引 4096 → 物理位置 376 8192 → 物理位置 782 每 4KB 数据写一条(log.index.interval.bytes) 二分② .log 物理位置 376 FileRecords.search 再顺序扫描几条即到 5000 索引即 mmap segment 滚动(写新段) log.segment.bytes = 1GB(默认) log.segment.ms = 7 天(默认) 只有活跃段可写 → 保证单文件顺序追加 两份索引(均 mmap 进虚拟内存) .index : offset → 物理位置 .timeindex: 时间戳 → offset 索引文件上限 10MB(log.index.size.max.bytes)

读任意 offset = 二分 segment + 二分索引 + 少量顺序扫描 —— 全程 O(1),与分区总数据量无关 写只需追加活跃段;删除整段回收;副本复制、消费读取复用同一套文件结构
存储三层:分区是目录,segment 是文件(1GB 或 7 天滚动,文件名是首条 offset),索引是小文件。关键设计是稀疏索引:不是每条消息一条索引,而是每 4KB 一条,所以 10MB 索引能覆盖 1GB 段。读取两段二分:先二分 segment 列表,再二分 mmap 索引,最后顺序扫几条消息。追问点:为什么索引用 mmap 而日志不用——索引小且大小固定(上限 10MB),映射一次即可;日志无限增长,mmap 频繁重映射得不偿失,顺序写 FileChannel 已够快。

Write Path

写路径:批量进、顺序追加、不主动 fsync

1 合流成批

同一分区的小消息在 producer 端聚成 batch;broker 把整批一次 append —— "把随机消息洪流变成线性写"

2 顺序追加

只写活跃 segment 末尾:FileRecords.append → FileChannel.write(源码 trunk),内核 write-behind 合并成大块物理写

3 进页缓存即算成功

log.flush.interval.messages 默认 Long.MAX_VALUE:从不主动 fsync,靠 OS 异步刷盘 + 副本冗余兜底

// trunk: clients/.../record/internal/FileRecords.java —— 追加只发生在文件末尾
public int append(MemoryRecords records) throws IOException {
    int written = records.writeFullyTo(channel);   // FileChannel 顺序写 → page cache
    size.getAndAdd(written);
    return written;
}

为什么敢不 fsync

官方设计文档:为一致性要求每次写都 fsync 会让性能下降两三个数量级。Kafka 的可靠性由复制协议保证——崩溃副本回归 ISR 前必须全量重新同步,丢没落盘的数据无所谓,日志以 Leader 为准截断对齐。

什么情况才配 fsync

只有单副本、或极端吞吐换持久化场景才调 log.flush.interval.messages / log.flush.interval.ms;每 60s 有 offset 检查点(log.flush.offset.checkpoint.interval.ms)用于重启恢复。

写路径三步:producer 聚批、broker 顺序 append、进页缓存即返回。高频追问"不 fsync 崩了不丢吗"——答案是不靠 fsync 靠副本:acks=all 时 ISR 全部(各自在内存/页缓存里)都有这条消息,一台宕机数据仍在其他机器的日志里;fsync 只解决单机持久化,Kafka 用复制把单机持久化问题变成了多机冗余问题。这一页可以和 MySQL redo log 对比:MySQL 是 WAL(先日志后数据页,还有随机写),Kafka 日志就是最终形态,全程顺序。

Zero-Copy

零拷贝:sendfile 让数据绕过用户态

传统 read/write 与 sendfile 对比 传统路径把数据从磁盘经页缓存拷到用户缓冲区再拷回内核 socket 缓冲,共四次拷贝两次系统调用;sendfile 路径数据从页缓存直接送网卡,共两次 DMA 拷贝一次系统调用,全程不进用户态。 传统 read() + write() —— 4 次拷贝 · 2 次系统调用 · 4 次上下文切换 用户空间 · 应用缓冲区 KERNEL Page Cache Socket 缓冲 磁盘 网卡 NIC ① DMA ② CPU ③ CPU ④ DMA read() write() sendfile() —— 2 次拷贝 · 1 次系统调用 · 2 次上下文切换 数据全程不进入用户态,页缓存中的数据被多个消费者反复复用 KERNEL Page Cache 磁盘 网卡 NIC ① DMA ② SG-DMA FileRecords.writeTo() → TransferableChannel.transferFrom() → FileChannel.transferTo()(sendfile) 失效场景:① SSL —— SslTransportLayer 读入 32KB direct buffer 交 SSLEngine 加密(源码实证)② 消息格式降级转换 down-conversion(KIP-110)
零拷贝是消费路径的大杀器。左边传统路径四拷贝:磁盘到页缓存(DMA)、页缓存到用户态(CPU)、用户态回 socket 缓冲(CPU)、socket 到网卡(DMA),外加两次系统调用四次上下文切换。sendfile 把中间两步 CPU 拷贝砍掉,网卡支持 scatter-gather 时连 socket 缓冲都不用,只剩两次 DMA。追问"什么时候失效"必须会:开 SSL 后 SslTransportLayer 用 32KB direct buffer 读取再加密,sendfile 作废(源码可证);消费端版本旧触发 down-conversion 时 broker 要读入内存转格式,同样失效。

Network Layer

网络层:Reactor 多路复用 + 两级线程池

1 Acceptor ×1

单个 acceptor 线程监听端口,把新连接轮询分发给 processor —— "a single acceptor thread and N processor threads"(官方 network-layer 文档)

2 Processor ×N

NIO 多路复用收发字节流,默认 num.network.threads=3(源码 SocketServerConfigs);读完整请求放入请求队列

3 Handler ×M

请求处理线程(KafkaApis)做真活:读写日志、更新 HW,默认 num.io.threads=8;完成后响应回 processor

为什么这样切分

  • 网络 IO 与业务逻辑解耦:慢请求(如磁盘 IO、追副本)不阻塞其他连接的收发
  • 请求队列 + 响应队列限流解耦,processor 用 Selector 管理成百上千连接

配套优化

  • 增量 Fetch Session(KIP-227):fetch 请求只带变化部分,空轮询开销骤降
  • socket 缓冲默认 100KB(socket.send.buffer.bytes),跨机房建议调大
  • 配额体系按 (user, client-id) 限带宽与线程占比,防单客户端打满 broker
网络层是经典 Reactor:一个 acceptor 负责建连,N 个 processor 负责 NIO 收发,M 个 io 线程执行请求。默认 3 个网络线程、8 个 IO 线程都来自源码。面试可展开:网络线程数对应连接吞吐,IO 线程对应磁盘吞吐;如果网络线程打满说明客户端请求太碎(要加批量),IO 线程打满说明磁盘或请求本身重。增量 fetch session 是消费者方向的补充优化。

Producer

生产者:RecordAccumulator 把小消息变成大请求

生产者批量发送流程 业务线程调用 send 后消息按分区进入 RecordAccumulator 的批队列,攒满 batch.size 或等待 linger.ms 超时后由 Sender 线程一次发出,单请求携带多个分区的 batch;压缩以整个 batch 为单位端到端保持。 APPEND 满/超时 PRODUCE 业务线程 send() 序列化 · 选分区 RECORDACCUMULATOR · buffer.memory = 32MB P0 batch P0 batch P0 … P1 batch P1 … P2 … 按分区聚合 · 同分区消息才可同批(保序单位) Sender 线程 linger.ms = 5ms Broker 整批顺序追加 端到端压缩(batch 级) gzip · snappy · lz4 · zstd producer 压整批 → broker 校验后原样落盘 → consumer 解压;批越大压缩比越高 默认值速查(trunk 源码) batch.size=16KB linger.ms=5ms buffer.memory=32MB max.request.size=1MB linger.ms 默认 0→5 是 4.0 的变化(效率提升且延迟反而更优) acks 与幂等 acks=all + enable.idempotence=true(3.0+ 默认) 幂等:PID + epoch + seq 去重重试,不产生重复 显式改 acks=1/0 会连带关闭幂等
生产者高吞吐的核心是 RecordAccumulator:send() 只是把消息按分区丢进批队列立即返回,Sender 线程在批满 16KB 或等满 5ms 时统一发出,一次 PRODUCE 请求可携带多个分区的批。三个易错点:① linger.ms 默认 5 是 4.0 起改的,旧资料写 0;② 压缩单位是整个 batch,且是端到端的——broker 不解压重压,原样落盘原样转发;③ 改 acks 为非 all 会静默关闭幂等。调优方向:追求吞吐加大 batch.size/linger.ms,追求延迟则减小。

Message Format

RecordBatch v2(magic=2):一切优化的载体

// 官方 message-format 文档(0.11+,当前 magic=2)
baseOffset: int64          // 批内第一条的 offset
batchLength: int32
partitionLeaderEpoch: int32
magic: int8  // 2
crc: uint32  // CRC-32C(硬件加速)
attributes: int16  // bit0-2: 压缩码 0无/1gzip/2snappy/3lz4/4zstd
                   // bit3: 时间戳类型 bit4: 事务 bit5: 控制批
lastOffsetDelta: int32
baseTimestamp / maxTimestamp: int64
producerId: int64  producerEpoch: int16
baseSequence: int32  // 幂等去重
recordsCount: int32
records: [ varint length · offsetDelta · timestampDelta
           · key · value · headers ]

为什么统一格式至关重要

producer、broker、consumer 三端共用同一二进制格式 → 数据块可不经过反序列化直接转发与落盘,这是零拷贝和"broker 尽量不解压"的前提。

增量编码省空间

批内每条记录只存 offsetDelta / timestampDelta(varint,同 Protobuf),时间戳、位移的冗余被摊平;批头还省掉了逐条 CRC。

batch 是最小业务单位

网络传输、日志追加、副本复制、压缩、事务(producerId/epoch/baseSequence)全部以 batch 为单位;批越大,头开销占比越低。

消息格式页回答"优化落在哪里":统一二进制格式让 broker 变成哑管道——只校验不解码,收到什么存什么发什么。v2 的三个看点:CRC32C 硬件加速且只算批头之后;varint 增量编码把逐条元数据压到最小;事务与幂等字段(producerId、epoch、baseSequence)内建在批头。追问 down-conversion:老 consumer 会迫使 broker 内存中转格式,零拷贝失效,所以保持客户端版本统一也是性能项。

Consumer

消费者:Pull + 长轮询 + offset 即状态

为什么 pull 不 push

push 时 broker 掌握节奏,消费者慢就被打爆;pull 让消费者按自己能力拉取,落后了就追,天然契合批量(官方 Push vs. pull 一节)。

长轮询防空转

naive pull 会忙等:fetch.min.bytes=1 + fetch.max.wait.ms=500 —— broker 无数据时挂起请求最多 500ms,有足量数据立即返回。

消费状态 = 一个整数

分区在组内被独占消费,位置就是下一个 offset——"just one number for each partition",周期性 checkpoint,ACK 机制因此变得极便宜。

拉取量参数(源码默认值)

  • max.partition.fetch.bytes = 1MB 每分区单次拉取上限
  • fetch.max.bytes = 50MB 消费者单次请求总上限;broker 侧同名配置默认 55MB(同名不同端,面试勿混淆)
  • max.poll.records = 500 单次 poll 返回条数(只影响应用侧,不影响底层 fetch)

页缓存×零拷贝的组合收益

官方:"on a cluster where the consumers are mostly caught up you will see no read activity on the disks whatsoever" —— 消费者追平进度时数据只在页缓存里,sendfile 直接从缓存送到网卡,多消费者重复读同一段日志只拷贝一次。

消费者页三件事:pull 的理由(速率自控+天然批量)、长轮询的两个参数、消费状态极简(一个 offset 整数,不需要 broker 记每条消息的 ack 状态——传统 MQ 的 sent/consumed 两态及其丢失/重复难题在这里根本不存在)。拉取量参数给源码默认值。最后强调组合拳:消费者追平时磁盘零读 + 零拷贝直送网卡,吞吐逼近网络极限。

Replication

ISR:用"动态副本集"换吞吐

ISR 与高水位 Leader 日志有 10 条,Follower1 已复制 9 条、Follower2 已复制 10 条;高水位取 ISR 中最小 LEO 为 9,消费者只能读到 offset 8;Follower 通过 fetch 拉取复制,30 秒内没跟上会被移出 ISR。 PRODUCE P0 Leader(Broker 1) 0 1 2 3 4 5 6 7 8 9 HW = 9(消费者只读 < 9) LEO = 10(下一条写入位置) 虚框 = 已写未提交(第 10 条) P0 Follower(Broker 2) LEO = 9 · 已追平 HW P0 Follower(Broker 3) LEO = 10 · 已复制第 10 条 ISR 入选条件 ① 与 controller 会话活跃(KRaft 心跳) ② 30s 内追上末端(replica.lag.time.max.ms) 不满足即被移出,恢复后须全量重同步再回 提交与可见性 committed = ISR 全部写入(acks=all 语义) 消费者只读 < HW,绝不见未提交数据 min.insync.replicas 兜底 ISR 过小 全挂时默认等 ISR(unclean 选举 = false)
ISR 是 Kafka 复制的灵魂:不搞 Raft 式多数派,而是维护一个动态追平集合,f+1 副本容 f 故障,2 副本就能容 1 台宕机,空间与写放大都省一半。代价是 acks=all 时提交要等全部 ISR(包括最慢的),多数派的"延迟只取决于最快副本"优势 Kafka 放弃了——官方认为允许客户端自选等待级别后,这个代价可接受。图里重点讲两个数字:LEO 是各副本日志末端,HW 是 ISR 最小 LEO,消费者永远读不到 HW 之后;第 10 条只有 Leader 有,就是典型的未提交数据。追不上 30 秒踢出 ISR,回来要全量重同步,所以不依赖 fsync 也不丢已提交数据。

Parallelism

分区:并行的度量衡,也是代价的来源

分区为什么带来吞吐

  • 顺序写的并行单位:每分区独立文件独立追加 → 写吞吐 ≈ Σ(磁盘带宽)
  • Leader 散布多 broker:生产/消费压力天然负载均衡
  • 消费组并行度上限 = 分区数(一个分区组内只归一个消费者)

分区过多的代价

  • 每分区每副本独立文件句柄与内存(producer 按 batch.size 预分配缓冲)
  • controller 元数据与复制协议压力、故障切换要迁更多 Leader
  • rebalance 更慢;端到端延迟上升(批更难攒满 + 副本间 fetch 更碎)

经验公式

分区数 ≥ max(目标生产吞吐 ÷ 单分区生产吞吐, 目标消费并行度),再按 broker 数留余量

先测单分区能力

用 kafka-producer-perf-test 实测单分区吞吐再外推,不要拍脑袋"越多越好"

扩分区不可逆

key→分区映射改变会破坏同 key 保序与局部性 —— 规划期定好,扩容提前

分区是把双刃剑。正面:它是顺序写和消费并行度的度量衡,吞吐随分区近似线性涨。反面:每个分区在每个副本上都是一组文件句柄和一份内存缓冲,几千个分区时元数据、rebalance、故障切换全是压力,批也更难攒满导致延迟上升。给出方法论:实测单分区吞吐、按公式定分区数、尽量在规划期定死(扩分区会打破 key 分区映射)。

Tuning Cheatsheet

高吞吐调优速查表(默认值均出自 trunk 源码)

Producer默认吞吐向调法
batch.size16 KB调到 64–512KB,攒更大批
linger.ms5 ms(4.0 前 0)50–100ms,用延迟换批大小
compression.typenonelz4 / zstd(zstd 压缩比更高,2.1+)
buffer.memory32 MB大批量场景加大,防 send() 阻塞
acks / 幂等all / true降为 1 可提吞吐但放弃可靠性与幂等
Broker / Topic默认吞吐向调法
num.network.threads3CPU 富余时上调(收发吞吐)
num.io.threads8磁盘是瓶颈时上调
log.segment.bytes1 GB过大拖慢删除与恢复,保持默认居多
log.flush.*不主动 fsync保持默认,靠副本保可靠
socket.send.buffer.bytes100 KB跨机房/高带宽延迟积场景调大
Consumer默认吞吐向调法
fetch.min.bytes / max.wait.ms1 / 500ms调大 min.bytes 让 broker 攒批响应
max.partition.fetch.bytes1 MB配合消息大小与处理能力上调
max.poll.records500单次处理更多条,摊薄 poll 开销

调优先看瓶颈在哪一段

  • 客户端 CPU 高 → 换更快压缩或减压缩等级;网络高 → 提压缩比
  • broker 网络线程打满 → 请求太碎,先加批;IO 线程打满 → 磁盘或请求重
  • 消费者 lag 大 → 先加分区/消费者,再查单条处理耗时
速查表按 producer、broker、consumer 三张表组织,全部默认值来自源码。讲的时候强调方法论:先定位瓶颈段(客户端 CPU / 网络 / 磁盘 / 消费处理),再对症调参。最常用的三个旋钮:batch.size、linger.ms、compression.type;消费者侧 fetch.min.bytes 与 max.poll.records。可靠性相关(acks、min.insync.replicas)动之前想清楚业务代价。

Cheat Sheet · 1/2

一页带走(上):六大支柱与必背数字

① Kafka 为什么快:六大支柱

顺序追加写只写活跃段末尾;顺序 600MB/s vs 随机 100kB/s,差 6000×
OS 页缓存读写全走堆外 page cache,零 GC;重启仍温热;32GB 机白得 28–30GB 缓存
sendfile 零拷贝FileRecords.writeTo → transferTo,数据不过用户态,4 次拷贝降到 2 次
批量 + 端到端压缩batch 是传输/存储/压缩共同单位;producer 压 → broker 原样存 → consumer 解
分区并行并行度与顺序写的交换单元;Leader 散布多机多盘,吞吐近线性
稀疏索引 + 二分每 4KB 一条 mmap 索引;读任意 offset = 二分段 + 二分索引 + 少量顺扫,全 O(1)

② 必背数字(默认值出自 trunk 源码)

segment 滚动log.segment.bytes=1GB · log.segment.ms=7天 · 索引上限 10MB
索引密度4KB 一条索引(log.index.interval.bytes
生产者批batch.size=16KB · linger.ms=5ms(4.0 前是 0)· buffer.memory=32MB
网络 / IO 线程num.network.threads=3 · num.io.threads=8 · socket 缓冲 100KB
消费者拉取fetch.min.bytes=1 + fetch.max.wait.ms=500 · max.partition.fetch.bytes=1MB · max.poll.records=500
副本活性replica.lag.time.max.ms=30s 追不上即被移出 ISR
速查上页:左栏是答"Kafka 为什么快"的总纲(六大支柱,每条一句话 + 一个数字),右栏是必背的默认值与量级。两栏都是回看与考前扫一眼用,细节回对应页码看推导。

Cheat Sheet · 2/2

一页带走(下):高频易错点与横向对比

③ 高频易错点(面试常在这翻车)

不 fsync 会丢吗fsync 只保单机;可靠性靠副本冗余,崩溃副本回归 ISR 前必须全量重同步
索引 mmap、日志不 mmap索引小且上限 10MB,映射一次复用;日志无限增长,顺序 FileChannel 已够快
ISR vs 多数派ISR f+1 容 f,比多数派 2f+1 省一半磁盘;代价是 acks=all 受最慢成员拖累
分区不是越多越好句柄/内存/元数据/rebalance 成本上升,批更难攒满 → 延迟上升
零拷贝会失效SSL(读入 32KB direct buffer 加密)或触发 down-conversion 时作废
同名配置两端不同fetch.max.bytes:客户端 50MB · broker 侧 55MB(同名不同端)

④ 与同类中间件:持久化差异

Kafka分区多文件 · FileChannel 顺序写 · sendfile 零拷贝
RocketMQ单 CommitLog 混合存储 · mmap 读写(单文件映射,省段管理)
PulsarBookKeeper 分层 · journal 顺序写 · 存算分离

⑤ 答题模板:性能题怎么开口

结构 → 支柱 → 数字 → 取舍。
· 结构:追加日志让一切读写变线性,性能与数据量解耦
· 支柱:顺序写 / 页缓存 / 零拷贝 / 批量压缩 / 分区并行 / 稀疏索引
· 数字:600MB/s vs 100kB/s(6000×)、4 次拷贝→2 次、28–30GB/32GB 页缓存
· 取舍:分区不是越多越好;acks 降级会静默关闭幂等;可靠性靠副本不靠 fsync
速查下页:六条最容易答错/被追问的点,一张三种中间件的持久化差异表(对比题先讲共性再讲差异),最后给"性能题怎么开口"的四步模板——结构、支柱、数字、取舍,缺一项都容易被判成背答案。

Interview QA · 1/2

典型面试 QA(上)

先盖住答案自己答一遍,再展开对照——想不起来比看得顺眼记得牢

Q1 Kafka 为什么吞吐这么高?

顺序写page cache零拷贝批量压缩分区并行

答六大支柱并给数字:顺序写 600MB/s vs 随机写 100kB/s(6000×);堆外页缓存 28–30GB/32GB 机;sendfile 4 拷贝→2 拷贝;批是网络/磁盘/压缩共同单位;分区水平扩展。收尾:结构决定性能——追加日志让一切变线性。

Q2 零拷贝原理?什么时候失效?

sendfiletransferToSSLdown-conversion

用户态中转的 4 次拷贝被砍成 2 次 DMA;源码链 FileRecords.writeTo → transferFrom → FileChannel.transferTo。失效:开启 SSL(SslTransportLayer 读 32KB direct buffer 加密)与消息格式降级转换(KIP-110)。

Q3 为什么索引用 mmap、日志却不用?

AbstractIndexFileChannel文件寿命

索引文件小且尺寸有上限(10MB),映射一次长期复用,随机访问收益大;日志无限增长且纯顺序写,mmap 要频繁重映射、缺页不可控,FileChannel 顺序写 + write-behind 已是线性速度。

Q4 写入不 fsync,崩溃会丢数据吗?

页缓存副本冗余ISR 重同步

fsync 只保单机;Kafka 用复制保可靠——acks=all 时消息已到全部 ISR(各自页缓存),单机断电数据仍在其他节点;崩溃副本回归前必须全量重同步。要更强持久化才配 log.flush.*,代价 2–3 个数量级。

Q5 页缓存为什么优于自建堆内缓存?

双份存储GC冷启动

堆内缓存会与 pagecache 双份存数据;对象开销与 GC 使大堆恶化;重启后堆内冷缓存而页缓存温热;一致性逻辑交给内核。32GB 机器等于白得 28–30GB"缓存"。

Q6 压缩算法怎么选?压缩发生在哪?

batch 级端到端lz4/zstd

压缩在 producer 侧按整批进行(同批冗余多、压缩比高),broker 校验后原样存储转发,consumer 解压——全链路只压一次。选型:lz4 均衡默认推荐,zstd 压缩比最高(2.1+),gzip 兼容好但 CPU 贵,snappy 通用。

Q7 分区越多吞吐越高吗?

边际递减fd/内存rebalance

不是。初期近线性;过多后文件句柄、producer 分区缓冲、controller 元数据、故障切换与 rebalance 成本上升,攒批更难、延迟更高。按"实测单分区吞吐 + 消费并行度"定分区数。

Q8 为什么用 offset 而不是 GUID 定位消息?

稀疏索引O(1)

GUID 需要维护随机 ID→offset 的全量映射,等于持久化随机索引;offset 单调递增,天然可二分稀疏索引,且分区+offset 唯一定位消息,索引结构极简——这是存储模型 O(1) 的根基(官方 implementation/log 文档)。

QA 上半场 8 题,覆盖吞吐主轴。每题答法:结构原因→数字→权衡。Q2 的源码证据链和 Q3 的 mmap 辨析是最容易拉开差距的两题;Q4 注意区分"fsync 保单机、复制保集群";Q7 要主动说"不是",展示对代价的理解。

Interview QA · 2/2

典型面试 QA(下)

同样先自答:这页偏对比与机制辨析,答完补一句"代价/边界"才完整。

Q9 为什么 O(1) 日志结构优于 BTree?

寻道成本超线性衰减

BTree O(log N) 在磁盘上意味着多次寻道,且磁盘串行寻道难并行;缓存混合场景下性能随数据量超线性衰减。追加日志读写全 O(1),性能与数据量解耦,因此敢于按时间而非按消费删除数据。

Q10 消费者如何避免空轮询打爆 broker?

长轮询fetch session

fetch.min.bytes(默认 1B)+ fetch.max.wait.ms(500ms):无数据时 broker 挂住请求直到凑够或超时;增量 fetch session(KIP-227)让重复请求只带增量分区集,减小空响应成本。

Q11 ISR 机制与 Raft 多数派的取舍?

f+1 容 f空间减半受最慢拖累

ISR:f+1 副本容 f 故障,比多数派(2f+1)省一半磁盘与写放大,Leader 候选全部来自 ISR;代价是 acks=all 的提交延迟取决于最慢 ISR 成员(多数派只需最快过半)。官方判断:可自选 acks 级别后,省副本更划算。

Q12 LEO 与 HW 是什么?消费者能读到什么?

LEOHW只读<HW

LEO 是副本日志下一条写入位置;HW = ISR 中最小 LEO。只有 ISR 全部写入(提交)的数据才推进 HW,消费者最多读到 HW-1,因此永远看不到可能丢的未提交数据;acks=all + min.insync.replicas 共同保证提交质量。

Q13 幂等生产者与事务的区别?

PID+seq跨分区原子

幂等:单分区、单会话内用 PID+epoch+序列号去重,重试不重复(默认开启);事务:跨多分区原子写入 + 与消费位点提交绑定(__consumer_offsets 同事务),配合 read_committed 实现流处理 exactly-once。

Q14 Kafka / RocketMQ / Pulsar 持久化差异?

sendfilemmap存算分离

Kafka:分区多文件 + FileChannel 顺序写 + sendfile 零拷贝;RocketMQ:单 CommitLog 混合存储 + mmap 读写(单文件映射,省去段管理);Pulsar:BookKeeper 分层,journal 顺序写 + 存算分离。零拷贝选型不同是常见追问点。

Q15 KRaft 是什么?为什么取代 ZooKeeper?

元数据即日志4.0 移除 ZK

把集群元数据(broker 注册、分区分配、配置)做成内部 Raft 日志 topic,由 controller quorum 复制;消除"两套一致性系统"(ZK + Kafka 自身 ISR)的运维与一致性负担,故障切换更快,支撑更大分区规模。4.0 起彻底移除 ZK。

下半场 7 题,偏向对比与机制辨析。Q11 和 Q12 是复制协议的核心,建议配合第 12 页的图讲;Q13 区分幂等(单分区会话)与事务(跨分区+位点原子);Q14 对比题先共同点后差异;Q15 用"元数据也是日志"一句话立住 KRaft 的本质。

Related & References

相关知识点与参考

本库相关 deck

  • TCP 连接管理 —— 零拷贝最终仍走 socket 发送,理解 sendfile 前先懂 TCP 缓冲与挥手
  • Go GC —— Kafka "堆外 page cache 优于堆内缓存" 的反面教材:大堆对 GC 的影响
  • MySQL MVCC —— 另一种"版本链 + 只读已提交"的思路,与 HW 语义对照
  • OS · 零拷贝 —— sendfile / SG-DMA / splice 的内核机制详解,本 deck 消费路径的理论底座

本领域深度 deck(已沉淀)

参考链接(一手来源)

收尾页给出本库内可跳转的相关 deck 与三个后续沉淀方向(可靠性矩阵、消费者组与 Rebalance、事务实现)。参考链接全部为官方文档、官方仓库源码与原始论文/文章,面试被深挖时可溯源。