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/s 6×7200rpm SATA RAID-5 顺序写(官方文档)
100 kB/s 同盘阵随机写 —— 差距 6000×
28–30 GB 32GB 机器上可用的 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 都是常数级
六大支柱是回答"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.size 16 KB 调到 64–512KB,攒更大批
linger.ms 5 ms(4.0 前 0) 50–100ms,用延迟换批大小
compression.type none lz4 / zstd(zstd 压缩比更高,2.1+)
buffer.memory 32 MB 大批量场景加大,防 send() 阻塞
acks / 幂等 all / true 降为 1 可提吞吐但放弃可靠性与幂等
Broker / Topic 默认 吞吐向调法
num.network.threads 3 CPU 富余时上调(收发吞吐)
num.io.threads 8 磁盘是瓶颈时上调
log.segment.bytes 1 GB 过大拖慢删除与恢复,保持默认居多
log.flush.* 不主动 fsync 保持默认,靠副本保可靠
socket.send.buffer.bytes 100 KB 跨机房/高带宽延迟积场景调大
Consumer 默认 吞吐向调法
fetch.min.bytes / max.wait.ms 1 / 500ms 调大 min.bytes 让 broker 攒批响应
max.partition.fetch.bytes 1 MB 配合消息大小与处理能力上调
max.poll.records 500 单次处理更多条,摊薄 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 读写(单文件映射,省段管理)
Pulsar BookKeeper 分层 · 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 零拷贝原理?什么时候失效?
sendfile transferTo SSL down-conversion
用户态中转的 4 次拷贝被砍成 2 次 DMA;源码链 FileRecords.writeTo → transferFrom → FileChannel.transferTo。失效:开启 SSL(SslTransportLayer 读 32KB direct buffer 加密)与消息格式降级转换(KIP-110)。
Q3 为什么索引用 mmap、日志却不用?
AbstractIndex FileChannel 文件寿命
索引文件小且尺寸有上限(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 是什么?消费者能读到什么?
LEO HW 只读<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 持久化差异?
sendfile mmap 存算分离
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(已沉淀)
生产者内核 —— acks / 重试与幂等 / 事务的完整可靠性矩阵
消费者与消费者组 —— 协调者、分配协议(Range/RoundRobin/Sticky/Cooperative)、Rebalance
副本机制与数据一致性(含 KRaft) —— ISR / HW / LeaderEpoch、unclean 选举、事务端到端
收尾页给出本库内可跳转的相关 deck 与三个后续沉淀方向(可靠性矩阵、消费者组与 Rebalance、事务实现)。参考链接全部为官方文档、官方仓库源码与原始论文/文章,面试被深挖时可溯源。