Skip to content

第 8 章 高吞吐的底层原理:Kafka 为什么这么快?

目标读者:知道「Kafka 很快」、面试常被问「Kafka 为什么快」却只能背出「顺序写、零拷贝、批量、压缩」四个名词然后被追问就答不上来的同学。

学完你会:能用画图的方式说清楚「sendfile 比 read+write 少了哪几次拷贝」、「PageCache 为什么帮 Kafka 把 Restart 的代价摊掉」、「mmap 为什么对 .index 文件特别合适」、「批量 + 压缩在 Producer/Broker/Consumer 三端是怎么协同的」。


0. 导读:Kafka 速度的「五大基石」

+---------------------------------------------------+
|                  Kafka 高吞吐五大基石               |
+---------------------------------------------------+
|  1. 磁盘顺序写  (Sequential Disk I/O)              |
|  2. PageCache    (内核页缓存替代 JVM 堆缓存)        |
|  3. 零拷贝       (sendfile / transferTo)           |
|  4. mmap         (索引文件直接映射内存)            |
|  5. 批量 + 压缩  (Producer/Broker/Consumer 协同)    |
+---------------------------------------------------+

关键认知:Kafka 几乎没有自己写过新的「优化算法」。它做的事情,是把 Linux 内核早就提供的特性 用对、用够、用满

  • 用对顺序写(磁盘最快的访问模式)。
  • 用对 PageCache(OS 已经为你做好的缓存层)。
  • 用对 sendfile(OS 早就提供的「数据从文件直送 socket」)。
  • 用对 mmap(让索引访问像访问数组一样)。
  • 用对批量 + 压缩(让 N 条消息共摊一次系统调用)。

理解了这五点,你就理解了 Kafka 为什么能做到「单台普通 Broker 几十万 msg/s、几百 MB/s」的吞吐。


1. 一个生活类比:快递分拣中心

把 Kafka 想象成一个超大快递分拣中心:

现实世界Kafka 对应
包裹按时间顺序往传送带上一直加顺序写 .log
中转区不开仓库存放,全部就地分装PageCache(不在 JVM 堆里复制一份)
装车时不开箱、按整辆车整体搬sendfile(数据不进用户态)
包裹号码在墙上贴二维码索引mmap 的 .index
按一车一车的批次而不是一件一件运batch + 压缩

反例:传统数据库就像「百货商店退换货」,每包裹都要拆开、入库、贴标签、再出库。慢的是这套「拆 + 装 + 拷」流程,不是快递本身。


2. 磁盘顺序写 vs 随机写

2.1 真实数量级

设备顺序写随机写 (4KB)倍数
HDD(7200 rpm SATA)~150 MB/s [典型值]~1 MB/s(约 100 IOPS)[典型值]150×
SATA SSD~500 MB/s [典型值]~50 MB/s(约 12K IOPS)[典型值]10×
NVMe SSD(消费级)~3000 MB/s [典型值]~600 MB/s(约 150K IOPS)[典型值]
NVMe SSD(企业级)~6000 MB/s [典型值]~1500 MB/s(约 350K IOPS)[典型值]

数据来源:Western Digital / Seagate / Samsung 公开 datasheet 中位数;实测会因队列深度、块大小不同而浮动。

结论:哪怕在 NVMe 上,顺序写也比随机写快 4~5 倍。HDD 上更是天上地下。

2.2 为什么差这么多

HDD 的物理瓶颈

  • 一次随机写需要:寻道(~5ms)+ 旋转延迟(~4ms,7200rpm)+ 数据传输(< 1ms)≈ 10ms / 4KB
  • 顺序写:没有寻道,没有旋转延迟,磁头一直跟着数据走 ≈ 0.04ms / 4KB
HDD 寻道时间分解(写一个 4KB 块):

随机写:    [寻道 5ms] [旋转 4ms] [传输 0.04ms]   = 9.04ms
顺序写:    [传输 0.04ms]                          = 0.04ms

倍数:9.04 / 0.04 ≈ 226×

SSD 的瓶颈虽小但仍存在

  • SSD 没有机械臂,但有 NAND 闪存的 Page / Block 结构
  • 每个 Block(典型 256 KB)只能整块擦除,写一个新 Block 前要先把 Block 内所有有效 Page 复制走(GC,垃圾回收)。
  • 顺序写:整 Block 写满、整 Block 擦除,几乎没有 GC 开销。
  • 随机写:Block 内有效 Page 散乱,GC 频繁,写放大 2x ~ 10x。

2.3 Kafka 怎么用

Kafka 的写入路径只有 「append 到 active segment 的 .log 末尾」 一种方式:

Producer.send()  →  Broker 接收  →  把 batch 字节追加到 .log 末尾

                                   纯顺序写
  • 绝不更新已有消息(offset 不可变)。
  • 绝不删除单条消息(要删就整个 segment unlink)。
  • 更新元数据用「另一条消息」(如 commit/abort marker、tombstone for compaction),仍然是 append-only。

这就是 Kafka 把 OS 顺序写的红利吃满的根本原因。

2.4 与数据库对比

InnoDB 也想顺序写(redo log 是顺序追加的),但 数据页随机修改

InnoDB 写入:
  1. 修改 buffer pool 中的数据页(内存随机写)
  2. 顺序写 redo log
  3. 异步刷脏页 → 磁盘随机写

Kafka 没有「数据页」概念,数据本身就是日志,写完就完事了。


3. PageCache:Broker 不在 JVM 堆里缓存消息

3.1 PageCache 是什么

PageCache 是 Linux 内核维护的「文件缓存」

  • 当你 read() 一个文件,内核先检查这个文件的某段是否在 PageCache。命中直接拷给用户态,没命中先从磁盘读进 PageCache 再拷。
  • 当你 write() 一个文件,数据先进 PageCache(标记 dirty),后台线程异步刷到磁盘。
  • PageCache 大小是所有可用空闲内存 - 应用占用,OS 自动管理。
                +---------------------------+
                |    Application (JVM)      |
                +-------------+-------------+
                              |
                              | read/write
                              v
                +---------------------------+
                |    PageCache (内核)        |   <- 同一份数据在内核里,所有进程共享
                +-------------+-------------+
                              |
                              | flush 异步
                              v
                +---------------------------+
                |        Disk (Hard)         |
                +---------------------------+

3.2 Kafka 为什么不自己在 JVM 堆里缓存

很多消息系统(如 RabbitMQ、ActiveMQ)会在自己的内存里维护一份消息副本,方便快速发给 Consumer。Kafka 故意不这么做

理由 1:JVM 对象开销大

JVM 中一个 byte[1024] 实际占用 ≈ 1024 + 16 (对象头) + 4 (length) + 8 (引用) = 1052 字节。如果要把每条消息包成 MessageRecord(包含 timestamp / key / value / headers),1KB payload 实际占用可能 1.5KB+。

理由 2:JVM 堆缓存会引发 GC

堆里几十 GB 的消息对象,每次 Major GC 暂停几秒到几十秒。Kafka 设计目标是永远不卡顿,所以宁可不缓存,让 OS 缓存。

理由 3:「同一份数据」在两处缓存是浪费

如果 JVM 堆和 PageCache 同时存有这条消息,等于在内存里存了两份。OS 可能会因为内存压力把 PageCache 里那份淘汰掉,等到下次读时又要从堆里复制一份过去——白白浪费。

理由 4:Broker 重启不丢缓存

JVM 进程重启 -> 堆里所有缓存数据丢失,几分钟才能预热完。 PageCache 在内核里 -> 进程重启不影响 PageCache,Broker 一起来就能享受热数据。

这一点是面试加分项:「Kafka Broker 滚动升级时几乎无吞吐损失」的核心原因。

3.3 一条消息的「写读路径」

写路径:

Producer  ->  Broker socket  ->  内存 buffer (短暂)  ->  write(.log)  ->  PageCache  ->  [异步] flush 到 disk

读路径(消费者拉的是热数据,比如刚写完几秒就消费):

Consumer  <-  socket (sendfile)  <-  PageCache (无需读盘)

热数据从来不下到磁盘——这就是 Kafka 在 PageCache 充足时能轻松达到 网卡上限 吞吐的根本。

3.4 配置建议

properties
# 不要设很大的 heap(让出空间给 PageCache)
KAFKA_HEAP_OPTS="-Xms6g -Xmx6g"   # 推荐 6~8GB,不要超过 16GB

# 不要主动 flush(让 OS 决定,靠副本机制保证可靠性)
log.flush.interval.messages=9223372036854775807
log.flush.interval.ms=9223372036854775807

反例:见过有人把 Kafka 设成 Xmx=64g,剩下不到 1GB 给 PageCache,结果 Broker 比磁盘还慢。

3.5 监控

bash
# 看 PageCache 命中率
sar -B 1
# %B/s = pgscank(被扫描)和 pgscand(被回收)越低越好

# 看 PageCache 占用
free -h
# buff/cache 这一栏就是 PageCache 大小

# 看具体某个文件多少在 PageCache
fincore /var/lib/kafka/data/learn.07.storage-0/00000000000000000000.log

健康的 Kafka Broker 应该有 70%+ 的 read 是 PageCache hit


4. 零拷贝:sendfile 比 read+write 少几次拷贝?

4.1 传统方式:4 次拷贝 + 2 次切换

「Broker 把磁盘上的消息发给 Consumer」这件事,传统做法(不用 sendfile)是这样的:

   Consumer 网络
        ^  socket
        |
 +------+------+
 |   Kernel    |   Stage 4: 内核 -> NIC(DMA)
 |   socket    |
 |   buffer    |   Stage 3: 用户态 -> 内核 socket buffer(CPU copy)
 |             |
 |   <bound>   |
 |             |
 |   PageCache |   Stage 1: 磁盘 -> PageCache(DMA)
 +------+------+
        |
        v
        Disk

详细 4 次拷贝 + 2 次上下文切换

1. 用户进程调 read() -> 上下文切换 [user -> kernel]
   
2. 内核:DMA 读磁盘 -> PageCache              [拷贝 1:磁盘 -> PageCache,DMA 不占 CPU]
   
3. 内核:CPU 把 PageCache 复制到 user buffer  [拷贝 2:PageCache -> 用户内存,CPU copy]
   上下文切换 [kernel -> user]
4. 用户进程拿到数据,调 write(socket)
   上下文切换 [user -> kernel]
5. 内核:CPU 把 user buffer 复制到 socket buffer  [拷贝 3:用户内存 -> socket buffer,CPU copy]
   
6. 内核:DMA 把 socket buffer 发到 NIC          [拷贝 4:socket buffer -> NIC,DMA 不占 CPU]
   上下文切换 [kernel -> user]

4 次拷贝(其中 2 次是 CPU 拷贝,吃 CPU)+ 2 次上下文切换 + 2 次系统调用。

4.2 零拷贝:sendfile 把 4 次降到 2 次

Linux 提供 sendfile(out_fd, in_fd, offset, count) 系统调用:

   Consumer 网络
        ^  socket
        |
 +------+------+
 |   Kernel    |   Stage 2: PageCache -> NIC(DMA gather,可选)
 |             |
 |             |
 |   PageCache |   Stage 1: 磁盘 -> PageCache(DMA)
 +------+------+
        |
        v
        Disk

sendfile 的工作流

1. 用户进程调 sendfile(socket, file, offset, count)
   上下文切换 [user -> kernel]
2. 内核:DMA 读磁盘 -> PageCache               [拷贝 1:磁盘 -> PageCache,DMA]
   
3. 内核:DMA 从 PageCache 直接传到 NIC          [拷贝 2:PageCache -> NIC,DMA gather]
   上下文切换 [kernel -> user]

只剩 2 次拷贝(都是 DMA,CPU 不参与)+ 1 次上下文切换 + 1 次系统调用。

严格地说,2.4 内核之前 sendfile 仍需要 1 次 CPU 拷贝(PageCache -> socket buffer),2.4 之后内核加入了 DMA gather 优化,可以让 NIC 直接从 PageCache 读取,做到「真·零 CPU 拷贝」。

4.3 拷贝次数对比表

方式DMA 拷贝CPU 拷贝总拷贝上下文切换系统调用
read + write22422
sendfile(pre 2.4)21311
sendfile + DMA gather(2.4+)20211

4.4 Kafka 在哪用了 sendfile?

只在「Consumer 拉消息」这条路径上

java
// FileChannel#transferTo() 在 Linux 下底层调 sendfile()
fileChannel.transferTo(position, count, socketChannel);

源码位置:org.apache.kafka.common.network.PlaintextTransportLayer#transferFrom

注意

  • Producer 写消息进 Broker -> 走的是普通 write(),不是 sendfile(Producer 端数据来自网络 buffer,不是文件)。
  • Consumer 读消息出 Broker -> 走 sendfile,从 .log 文件直接到 socket。
  • SSL 加密的连接 -> 不能用 sendfile(数据要在用户态加密),所以 SSL 会显著降低吞吐。

4.5 实测数据 [典型值]

同一台机器(Intel Xeon 8 核 / 64GB 内存 / NVMe SSD / 10GbE 网卡):

协议吞吐CPU 占用
PLAINTEXT (sendfile)~900 MB/s~30%
SSL (无 sendfile)~400 MB/s~80%

这就是为什么 SSL Kafka 集群通常需要更多 CPU 核数。


5. mmap:让 .index 像数组一样访问

5.1 mmap 是什么

mmap() 把文件直接映射到进程的虚拟地址空间。访问内存地址 = 访问文件,不需要 read/write 系统调用

   进程虚拟地址空间
   +-----------------+
   |   ...           |
   |   ...           |
   |   mmap region   |  <-+
   |   (.index)      |    |
   |   ...           |    |  page fault -> 内核加载该 page
   +-----------------+    |
                          v
   +-----------------+
   |   PageCache     |  <- 内核
   |   (.index 文件)  |
   +-----------------+
                          ^
                          |  DMA
   +-----------------+
   |     Disk        |
   +-----------------+

5.2 Kafka 怎么用

org.apache.kafka.common.utils.OffsetIndex 在打开 .index 文件时调用:

java
mmap = raf.getChannel().map(FileChannel.MapMode.READ_WRITE, 0, length);
  • 整个 .index 文件被映射到进程地址空间。
  • 二分查找时直接对 mmap 数组做指针运算,没有 read 系统调用
  • 写新索引项时直接 mmap.putInt(rel)mmap.putInt(pos)没有 write 系统调用

5.3 mmap 的优势

维度read/writemmap
数据流用户 buffer 与 PageCache 互相 copyPageCache(不复制到用户态)
系统调用次数每次 IO 都要只有第一次 page fault
适合场景一次性大块 IO小块、随机、反复访问

5.4 为什么 .log 不用 mmap?

.log 文件单个就 1GB,mmap 进 32 位地址空间装不下。即便是 64 位,mmap 1GB 文件后所有 PageCache 都被 lock 住,难以被 OS 灵活调度。所以 Kafka 只对 .index / .timeindex 用 mmap(这俩各 ≤ 10MB)。

5.5 一个微妙的坑

mmap 写入是异步的——只标记 dirty page,不立刻写盘。Broker 崩溃时 mmap 的最后几条索引项可能丢失,但没关系:

  • .index 是稀疏的辅助结构,重启时 Kafka 会用 .log 重建索引(按 4KB 顺序扫一遍)。
  • .log 才是真相之源,索引丢了不损坏数据。

6. 批量 + 压缩:三端协同

6.1 Producer 端:先攒后发

Producer 不是来一条发一条,而是攒成 batch 再发

  send(record_1) -> 进 RecordAccumulator (内存 buffer)
  send(record_2) -> 同 batch
  send(record_3) -> 同 batch
  ...
  达到下面任一条件 -> Sender 线程整 batch 发出去:
    +- batch.size 满了(默认 16KB)
    +- linger.ms 到了(默认 0,但生产环境推荐 5~50)
    +- flush() 强制发

类比:快递员凑够一车 100 个包裹再送,比每来一个包裹送一次效率高 100 倍。

6.2 Producer 端:压缩(在 batch 维度)

Producer 把整个 batch 的 records 段一起压缩:

properties
compression.type=lz4   # none / gzip / snappy / lz4 / zstd

5 种压缩对比 [典型值]

算法压缩率压缩速度解压速度CPU 占用推荐场景
none1.0x--0局域网内 + 已经是二进制
gzip5~7x老遗留集群
snappy2~3xKafka 默认很久
lz43~4x极快极快生产首选
zstd4~6x重度压缩需求(成本敏感)

关键认识:批量 + 压缩有「化学反应」:一个 batch 内多条消息往往字段名、字段格式高度重复,压缩率非常高。如果不批量直接压缩单条,压缩率几乎等于 1。

6.3 Broker 端:原样存储

Broker 收到压缩 batch 后完全不解压,直接 append 到 .log

  • 节省 Broker CPU。
  • 节省磁盘空间(按压缩后存)。
  • 节省网络出口带宽(Consumer 拿到的也是压缩 batch)。

这就是 §6.2 RecordBatch 的 attributes 字段会带 compresscodec 的原因。

6.4 Consumer 端:拿到压缩 batch 后才解压

Broker.fetch -> Consumer 拿到压缩 batch -> Consumer 端解压 -> 一条条交给业务

解压发生在Consumer 进程,不是 Broker。这把解压的 CPU 开销分散到 N 个 Consumer 上。

6.5 端到端数据流

Producer:  records 攒 batch -> 压缩 batch -> 网络
                                              v
Broker:    收到压缩 batch -> 直接 append .log -> PageCache
                                              v
Consumer:  sendfile 拿到压缩 batch -> 解压 -> 业务消费

全程不解压、不重新压缩 —— 这是 Kafka 高吞吐的关键之一。

6.6 配置示例

properties
# Producer 推荐
compression.type=lz4
batch.size=65536           # 64 KB
linger.ms=20               # 等 20ms 凑大 batch
buffer.memory=67108864     # 64 MB 总缓冲
acks=all
enable.idempotence=true

# Topic 端推荐
compression.type=producer  # 透传 Producer 压缩,不重压缩

7. 端到端延迟与吞吐的权衡:linger.ms 曲线

7.1 公式

吞吐 ~ batch_size * send_rate
延迟 ~ linger.ms + 网络往返 + Broker 处理

7.2 实测曲线 [典型值]

同一集群、同一 producer、不同 linger.ms

linger.msbatch 平均大小吞吐 (MB/s)p99 延迟 (ms)
01 KB805
512 KB2008
2032 KB38025
5060 KB48055
10065 KB510105
50065 KB(达到 batch.size 上限)520505

观察

  1. linger 从 0 -> 20,吞吐提升 5 倍,延迟仅增加 20ms。
  2. linger 从 20 -> 100,吞吐只再提升 30%,延迟却涨 5 倍。
  3. 甜点(sweet spot)通常在 5 ~ 30ms,业务能接受的延迟内取最大。

7.3 业务取舍

业务场景推荐 linger.ms
实时风控(要求 <= 100ms 端到端)5 ~ 10
一般日志聚合20 ~ 50
离线批量 ETL(容忍秒级)100 ~ 500
低吞吐高优先级(订单事件)0 ~ 5

8. 与 RabbitMQ / RocketMQ / Pulsar 对比

维度KafkaRabbitMQRocketMQPulsar
主要存储模式顺序写 + PageCache + sendfileErlang Mnesia + 内存为主CommitLog 顺序写 + ConsumeQueue 索引Bookkeeper Ledger 多副本写
零拷贝sendfilesendfilesendfile
Broker 是否在堆里缓存消息否(依赖 PageCache)是(内存 + 持久化)否(PageCache)
批量Producer 端攒 batch单条 ack 多单条发送、batch 较少Producer 攒 batch
压缩粒度整个 batch单条单条整个 entry
典型单 Broker 吞吐几十万 msg/s, 几百 MB/s几万 msg/s几十万 msg/s数十万 msg/s
延迟 p995ms ~ 50ms(取决 linger.ms)1ms ~ 10ms5ms ~ 30ms5ms ~ 50ms

总结:Kafka 的「快」是牺牲了端到端延迟、增加批次延迟换来的吞吐。如果业务真的需要单条 ms 级延迟,RabbitMQ / NATS 反而更合适。


9. 一张图总结

                    Kafka 高吞吐全链路图

  +------------------+
  |     Producer     |
  +--------+---------+
           |   1. 在 RecordAccumulator 里 攒 batch
           |   2. 整 batch 用 lz4/zstd 压缩
           |   3. 一次网络 send 多条
           v
  +------------------+
  |      Broker      |
  | 1. socket 收 batch
  | 2. 直接 append    | <-- 顺序写 .log
  |    .log 末尾      |
  | 3. 写 PageCache   | <-- 内核异步刷盘
  |    (不解压)       |
  | 4. Follower 拉   |
  |    batch 同步     | <-- ISR 副本机制
  +--------+---------+
           |   5. Consumer 来 fetch
           |   6. sendfile:PageCache -> NIC(无 CPU 拷贝)
           v
  +------------------+
  |     Consumer     |
  | 1. 拿到压缩 batch |
  | 2. 解压(在 Consumer 端,分散 CPU 压力)
  | 3. 一条条交业务   |
  +------------------+

  关键点:从 Producer 到 Consumer,消息的「字节序列」几乎不变,
         不在任何中间节点做无谓的拆包/重新打包。

10. 小结

Kafka 高吞吐五大基石
+- 顺序写
|  +- HDD 顺序 100MB/s vs 随机 1MB/s = 100x
|  +- NVMe 顺序 3000 vs 随机 600 = 5x
|  +- Kafka 只 append .log 末尾
+- PageCache
|  +- Broker 不在 JVM 堆里缓存消息
|  +- Restart 不丢缓存
|  +- 留给 PageCache 50%+ 内存
+- 零拷贝(sendfile)
|  +- 4 次拷贝 + 2 次切换 -> 2 次 DMA 拷贝 + 1 次切换
|  +- 只在 Consumer 拉消息路径上用
|  +- SSL 用不了 sendfile,吞吐显著下降
+- mmap
|  +- .index / .timeindex 直接映射内存
|  +- 二分查找无系统调用
|  +- .log 不 mmap(太大)
+- 批量 + 压缩
   +- Producer 攒 batch + lz4 压缩
   +- Broker 不解压直接 append
   +- Consumer 端解压(分散 CPU)

11. 面试高频题(7 题)

Q1:Kafka 为什么这么快?

考察点:综合理解高吞吐底层原理。

答案

  1. 顺序写:消息只 append 到 active segment 的 .log 末尾。HDD 顺序写比随机写快 100x,SSD 上也快 5x。
  2. PageCache:Broker 不在 JVM 堆里缓存消息,靠 OS 的 PageCache 缓存。优势:
    • 避免 GC 暂停。
    • Broker 重启不丢缓存(PageCache 在内核里)。
    • 同一份数据不在堆和 PageCache 里同时存两份。
  3. 零拷贝:Consumer 拉数据走 sendfile(在 Kafka 里通过 FileChannel#transferTo)。从 read+write 的 4 次拷贝 + 2 次上下文切换,降到 2 次 DMA 拷贝 + 1 次上下文切换。
  4. mmap.index.timeindex 用 mmap,二分查找是纯内存指针运算。
  5. 批量 + 压缩:Producer 攒 batch、整 batch 压缩,Broker 原样落盘和转发,Consumer 端才解压。压缩算法首选 lz4(吞吐最优,CPU 低)。
  6. 加分项
    • SSL 不能 sendfile,会让吞吐腰斩。
    • Producer 端 linger.ms=20 通常是吞吐/延迟最优拐点。
    • 给 PageCache 预留足够内存(建议 JVM heap <= 8GB)。

Q2:sendfile 怎么做到「零拷贝」?比传统 read+write 少了几次拷贝?

考察点:操作系统、Linux 系统调用。

答案

  1. 传统 read + write 一共 4 次拷贝、2 次上下文切换
    • 磁盘 -> PageCache(DMA)
    • PageCache -> 用户 buffer(CPU)
    • 用户 buffer -> socket buffer(CPU)
    • socket buffer -> NIC(DMA)
  2. sendfile(pre 2.4 内核)3 次拷贝、1 次上下文切换
    • 磁盘 -> PageCache(DMA)
    • PageCache -> socket buffer(CPU)
    • socket buffer -> NIC(DMA)
  3. sendfile + DMA gather(2.4+ 内核)2 次 DMA 拷贝、1 次上下文切换
    • 磁盘 -> PageCache(DMA)
    • PageCache -> NIC(DMA gather,NIC 直接从 PageCache 读)
  4. 「零拷贝」中的「零」是指零 CPU 拷贝,不是物理上没有拷贝。
  5. Kafka 里的应用
    • Consumer fetch 走 FileChannel#transferTo -> JVM -> 底层 sendfile。
    • Producer 写 Broker 走 sendfile(数据来自网络 buffer 不是文件)。
    • SSL 连接走不了 sendfile(数据要先到用户态加密)。
  6. 加分项:提到 Java 21 之前 transferTo 在 SSL 下会 fallback 到普通 read/write,21 之后引入 io_uring 有进一步优化。

Q3:Kafka Broker 为什么不在 JVM 堆里缓存消息?

考察点:PageCache、JVM、设计哲学。

答案

  1. JVM 对象开销大byte[1024] 实际占用约 1052 字节,包装成 Message 对象后接近 1.5KB。
  2. GC 是定时炸弹:堆里几十 GB 的对象,每次 Major GC 暂停几秒到几十秒,Kafka 的吞吐目标受不了。
  3. 数据双份浪费:同一份消息如果堆里和 PageCache 里都存,纯属浪费。OS 在内存压力下还得去淘汰一份,反而恶化。
  4. Broker 重启不丢缓存:JVM 堆缓存进程退出就没;PageCache 在内核里,进程重启后还在。Kafka 滚动升级时几乎无吞吐损失。
  5. OS 已经把缓存做得很好:PageCache 自动 LRU、自动预读、自动写回,Kafka 没必要重新发明。
  6. 配置佐证:官方建议 Heap 6~8GB、剩下都让出给 PageCache;log.flush.interval.messages / log.flush.interval.ms 默认都是无穷大,意思是「不用主动 flush,让 OS 决定」。
  7. 加分项:提到 Kafka 是 Scala 写的(JVM 体系)但消息从来不进 Scala 对象层,全程 ByteBuffer,本质就是「让 OS 帮我做一切」。

Q4:Kafka 顺序写真的比随机写快两个数量级吗?给真实数字。

考察点:对存储硬件的认知。

答案

  1. HDD(机械盘 7200 rpm SATA)[典型值]
    • 顺序写:~150 MB/s
    • 随机写 4KB:~1 MB/s(约 100 IOPS)
    • 倍数:~150x,因为机械臂寻道(5ms) + 旋转延迟(4ms) 占绝对主导。
  2. SATA SSD [典型值]
    • 顺序:~500 MB/s
    • 随机 4KB:~50 MB/s
    • 倍数:~10x。SSD 没有寻道,但 NAND 的 GC 写放大让随机写慢一截。
  3. NVMe SSD(消费级)[典型值]
    • 顺序:~3 GB/s
    • 随机 4KB:~600 MB/s
    • 倍数:~5x。NVMe 仍然偏好顺序,但差距收窄。
  4. Kafka 怎么用:永远 append 到 .log 末尾,永远不修改已有消息。删除按 segment 整体 unlink,不是按消息删。
  5. 加分项:提到 RDBMS 的 redo log 也是顺序写,但数据页随机写让总体仍偏慢;Kafka 没有「数据页」,写完就完事。

Q5:mmap 在 Kafka 哪里用?为什么 .log 不 mmap?

考察点:mmap 原理、Kafka 存储结构。

答案

  1. 应用位置:仅 .index.timeindex 文件用 mmap(org.apache.kafka.common.utils.OffsetIndex)。
  2. mmap 优势
    • 把文件直接映射到进程虚拟地址空间,访问像访问数组指针。
    • 二分查找时没有 read 系统调用,只有第一次 page fault。
    • 写新索引项也是直接写 mmap,不需要 write 系统调用。
  3. .log 为什么不 mmap
    • 单文件 1GB,mmap 进程地址空间不灵活(32 位下根本装不下,64 位也会让 PageCache 锁住一大片)。
    • .log 的访问模式是「顺序读写」,read/write + sendfile 已经很高效,mmap 收益不大。
    • mmap 写盘是异步的,崩溃可能丢最后几个 page,但 .index 可以从 .log 重建,.log 不能丢
  4. 崩溃恢复.index 丢了几条索引没关系,Broker 启动时会用 .log 重新构建索引。
  5. 加分项:提到 Java NIO 的 FileChannel.map() 是 mmap 入口;MappedByteBuffer 在 GC 后才能真正 unmap,所以 Broker 关闭时 .index 文件可能要等几秒才能被 OS 真正释放。

Q6:批量 + 压缩在 Producer / Broker / Consumer 三端怎么协同?

考察点:Kafka 数据流。

答案

  1. Producer 端
    • 多条消息进入 RecordAccumulator 攒 batch(按 partition 维度)。
    • 凑够 batch.size(默认 16KB)或 linger.ms 到时间,整 batch 用 lz4/zstd 压缩。
    • 一次网络 send 把整个压缩 batch 发给 Broker。
  2. Broker 端
    • 收到压缩 batch 后完全不解压,直接 append 到 .log 末尾。
    • 既省 Broker CPU,又省磁盘空间。
    • Topic 端建议设 compression.type=producer 透传 Producer 压缩。
  3. Follower 副本:拉到的也是压缩 batch,同样不解压直接 append。
  4. Consumer 端
    • 通过 sendfile 拿到压缩 batch(连压缩字节都不变)。
    • 在 Consumer 进程解压、拆 batch 成单条 record 交给业务。
  5. 化学反应:一个 batch 内消息往往字段名、格式重复,压缩率比单条压缩高得多(5~10x)。
  6. 配置推荐
    • compression.type=lz4(吞吐最优、CPU 低)
    • batch.size=65536(64KB)
    • linger.ms=20
  7. 加分项:提到 Broker 如果设 compression.type=lz4 且 Producer 是 gzip,Broker 会重新压缩(解压再压),吞吐和 CPU 都受损;所以默认 producer 模式最好。

Q7:linger.ms 调高会怎样?给一个实测对比。

考察点:吞吐 vs 延迟取舍。

答案

  1. 作用:Producer 在没凑够 batch.size 时,最多等 linger.ms 毫秒再发,期望期间能凑出更大的 batch。
  2. 公式
    • 吞吐 ~ batch 平均大小(更大 batch -> 每次系统调用 + 压缩 + 网络携带的消息更多)。
    • 延迟 ~ linger.ms + 网络 RTT + Broker 处理(5~10ms)。
  3. 实测对比 [典型值](同一集群、同一 producer):
    • linger=0, batch.size=16K -> 吞吐 80MB/s, p99 延迟 5ms
    • linger=5 -> 200MB/s, 8ms
    • linger=20 -> 380MB/s, 25ms
    • linger=50 -> 480MB/s, 55ms
    • linger=100 -> 510MB/s, 105ms
  4. 拐点:通常在 5~30ms。再大吞吐增长缓慢,延迟却线性涨。
  5. 业务推荐
    • 实时风控/订单 -> linger=5 以下。
    • 日志聚合 -> linger=20~50
    • 离线 ETL -> linger=100~500
  6. 加分项:提到 flush() 可以强制不等 linger 立刻发;batch.size 太小会让 linger 没意义(攒满就发了);与 enable.idempotence 配合时 max.in.flight.requests.per.connection 必须 <= 5。

本章配套:

  • 08_performance/demo.html:四次拷贝 vs 两次拷贝动画 + 顺序写 vs 随机写磁头动画 + PageCache 命中率柱状图 + 三参数实时计算。
  • 08_performance/code/sequential_vs_random.py:用 Python os.write 模拟顺序写 vs 随机写文件吞吐对比。
  • 08_performance/code/throughput_bench.py:完整 Producer -> Consumer 端到端压测脚本。

延伸阅读

  • 第 7 章 存储与日志格式 —— 「顺序写」「mmap」依赖的存储模型。
  • 第 4 章 Producer 深入 —— linger.msbatch.sizecompression.type 的更细致讨论。
  • 第 5 章 Consumer 深入 —— fetch.min.bytesfetch.max.wait.ms 是 Consumer 侧的「批量」对应物。
  • 第 9 章 副本机制与 ISR —— Follower 拉取也走 sendfile + 批量。

🎬 可视化演示

演示加载缓慢或样式异常?点此在新标签页打开 ↗

💻 示例代码

python
#!/usr/bin/env python3
"""
sequential_vs_random.py
=======================

用 Python `os.write` + `os.lseek` 直接对裸文件做 4KB 块的「顺序写」vs「随机写」吞吐对比,
直观感受为什么 Kafka 的「顺序追加日志」是性能制胜法宝。

测量逻辑:
  - 顺序写:在一个空文件里从头到尾追加 N 个 4KB 块。
  - 随机写:在一个**预先分配好** N×4KB 大小的文件里,每次写一个随机偏移的 4KB 块。
  - 都用 `os.fsync` 强制刷盘,避免 PageCache 假象。

输出指标:
  - 总写入字节数
  - 耗时
  - 吞吐(MB/s)
  - 顺序 / 随机 倍数比

依赖:标准库(os / time / random / argparse),不需要任何第三方包。

用法:
  python3 sequential_vs_random.py
  python3 sequential_vs_random.py --blocks 50000 --block-size 4096
  python3 sequential_vs_random.py --no-fsync          # 看 PageCache 加持下的吞吐
  python3 sequential_vs_random.py --workdir /tmp      # 换数据目录(建议放 SSD/HDD/NVMe 各跑一次)
"""

from __future__ import annotations

import argparse
import os
import random
import sys
import time
import shutil
from pathlib import Path


def write_sequential(path: Path, blocks: int, block_size: int, fsync: bool) -> float:
    """
    在一个空文件中从头到尾追加 blocks 个 block_size 字节。
    返回耗时(秒)。
    """
    payload = b"\xAB" * block_size
    fd = os.open(str(path), os.O_CREAT | os.O_WRONLY | os.O_TRUNC, 0o644)
    t0 = time.perf_counter()
    try:
        for _ in range(blocks):
            os.write(fd, payload)
        if fsync:
            os.fsync(fd)
    finally:
        os.close(fd)
    return time.perf_counter() - t0


def write_random(path: Path, blocks: int, block_size: int, fsync: bool, seed: int) -> float:
    """
    在一个预分配 blocks*block_size 字节的文件里随机偏移写 blocks 个块。
    返回耗时(秒)。
    """
    total = blocks * block_size
    # 预分配
    with open(path, "wb") as f:
        f.truncate(total)

    rng = random.Random(seed)
    # 预生成所有偏移(不计入计时)
    offsets = [rng.randrange(0, blocks) * block_size for _ in range(blocks)]
    payload = b"\xCD" * block_size

    fd = os.open(str(path), os.O_WRONLY)
    t0 = time.perf_counter()
    try:
        for off in offsets:
            os.lseek(fd, off, os.SEEK_SET)
            os.write(fd, payload)
        if fsync:
            os.fsync(fd)
    finally:
        os.close(fd)
    return time.perf_counter() - t0


def human_mb(b: int) -> str:
    return f"{b / 1024 / 1024:.2f} MB"


def detect_fs_type(path: Path) -> str:
    """尽力探测挂载点 + 文件系统类型,方便读者解读。"""
    try:
        with open("/proc/self/mounts") as f:
            mounts = f.read().splitlines()
    except Exception:
        return "?"
    parent = path.resolve().parent
    while parent != Path("/") and parent.exists():
        for line in mounts:
            parts = line.split()
            if len(parts) >= 3 and parts[1] == str(parent):
                return f"{parts[0]} ({parts[2]})"
        parent = parent.parent
    return "?"


def main() -> None:
    parser = argparse.ArgumentParser(description="顺序写 vs 随机写吞吐对比")
    parser.add_argument("--workdir", type=str, default="./_seq_rand_bench",
                        help="测试数据目录(运行后会被删除)")
    parser.add_argument("--blocks", type=int, default=10000,
                        help="块数,默认 10000 个")
    parser.add_argument("--block-size", type=int, default=4096,
                        help="每块字节数,默认 4096")
    parser.add_argument("--no-fsync", action="store_true",
                        help="不调 fsync(看 PageCache 加持下的吞吐)")
    parser.add_argument("--seed", type=int, default=42)
    args = parser.parse_args()

    workdir = Path(args.workdir)
    workdir.mkdir(parents=True, exist_ok=True)

    seq_file = workdir / "sequential.bin"
    rnd_file = workdir / "random.bin"
    total_bytes = args.blocks * args.block_size

    print("=" * 72)
    print(" 顺序写 vs 随机写 吞吐对比")
    print("-" * 72)
    print(f"  workdir   : {workdir.resolve()}")
    print(f"  fs        : {detect_fs_type(workdir)}")
    print(f"  blocks    : {args.blocks:,}")
    print(f"  block_size: {args.block_size} B")
    print(f"  total     : {human_mb(total_bytes)}")
    print(f"  fsync     : {'NO (依赖 PageCache)' if args.no_fsync else 'YES (强制刷盘)'}")
    print("=" * 72)

    print("\n[1/2] 顺序写…")
    seq_time = write_sequential(seq_file, args.blocks, args.block_size, fsync=not args.no_fsync)
    seq_mbps = total_bytes / seq_time / 1024 / 1024
    print(f"  耗时 {seq_time*1000:8.1f} ms   吞吐 {seq_mbps:7.1f} MB/s")

    print("\n[2/2] 随机写…")
    rnd_time = write_random(rnd_file, args.blocks, args.block_size,
                            fsync=not args.no_fsync, seed=args.seed)
    rnd_mbps = total_bytes / rnd_time / 1024 / 1024
    print(f"  耗时 {rnd_time*1000:8.1f} ms   吞吐 {rnd_mbps:7.1f} MB/s")

    ratio = seq_mbps / rnd_mbps if rnd_mbps > 0 else float("inf")

    print("\n" + "=" * 72)
    print(f"  顺序 / 随机 = {ratio:.2f}x")
    if ratio > 50:
        print("  ✅ 数量级差异(典型 HDD 行为)。Kafka 的顺序追加日志极大利用了这点。")
    elif ratio > 5:
        print("  🟡 倍数差异(典型 SSD 行为)。即使 SSD,顺序写也明显占优。")
    elif ratio > 1.5:
        print("  🟢 小幅领先(NVMe 或 PageCache 假象)。如果用了 --no-fsync,差异会被 OS 缓存掩盖。")
    else:
        print("  ⚠️  几乎无差异。可能 fsync 没生效,或 workdir 在 tmpfs 上。")
    print("=" * 72)

    # 清理
    try:
        shutil.rmtree(workdir)
        print(f"\n  已清理 {workdir}")
    except Exception as e:
        print(f"\n  ⚠️  清理 {workdir} 失败:{e}")

    print("\n实验提示:")
    print("  - 在 HDD 上跑预期看到 50~200x 顺序优势")
    print("  - 在 SATA SSD 上跑预期看到 5~15x")
    print("  - 在 NVMe SSD 上跑预期看到 3~6x")
    print("  - 加 --no-fsync 后,PageCache 会让随机写也变快,但生产环境必须 fsync")


if __name__ == "__main__":
    sys.exit(main())
python
#!/usr/bin/env python3
"""
throughput_bench.py
===================

完整的 Producer → Broker → Consumer 端到端吞吐压测脚本。

核心目的:把第 8 章里讲到的「批量大小 / 压缩算法 / linger.ms / acks」
四个杠杆都暴露成命令行参数,让你在自己的集群上**亲手感受**它们对吞吐 / 延迟的影响。

依赖:confluent-kafka 或 kafka-python。脚本会自动检测,二选一。
  pip install confluent-kafka
  # 或
  pip install kafka-python

用法:
  python3 throughput_bench.py --bootstrap localhost:9092 --topic learn.bench
  python3 throughput_bench.py --records 200000 --size 1024 --batch 65536 --linger 20 --compression lz4
  python3 throughput_bench.py --acks all --compression zstd --no-consume   # 只压 Producer

输出指标:
  - Producer: 发送速率(msg/s, MB/s)、平均延迟、p50/p95/p99 延迟
  - Consumer: 消费速率(msg/s, MB/s)、首条消息延迟(端到端)

提示:
  - 别在生产集群跑大流量压测。
  - 跑之前先用 kafka-topics.sh 创建好 Topic,分区数 ≥ Producer 并发 / Consumer 并发 的最大值。
"""

from __future__ import annotations

import argparse
import os
import statistics
import sys
import threading
import time
import uuid
from dataclasses import dataclass, field
from typing import List, Optional


# ----------------- 客户端封装:先尝试 confluent-kafka,回退 kafka-python -----------------

class KafkaImpl:
    name = "none"

    @classmethod
    def detect(cls) -> "KafkaImpl":
        try:
            import confluent_kafka  # noqa: F401
            return ConfluentImpl()
        except ImportError:
            pass
        try:
            import kafka  # noqa: F401
            return KafkaPythonImpl()
        except ImportError:
            pass
        raise RuntimeError(
            "未找到 confluent-kafka 或 kafka-python。请先:\n"
            "  pip install confluent-kafka   (推荐,性能更好)\n"
            "  或 pip install kafka-python"
        )


class ConfluentImpl(KafkaImpl):
    name = "confluent-kafka"

    def make_producer(self, args):
        from confluent_kafka import Producer
        cfg = {
            "bootstrap.servers": args.bootstrap,
            "acks": args.acks,
            "linger.ms": args.linger,
            "batch.size": args.batch,
            "compression.type": args.compression,
            "enable.idempotence": "true" if args.acks == "all" else "false",
            "queue.buffering.max.messages": 1_000_000,
        }
        return Producer(cfg)

    def produce(self, producer, topic, key, value, on_delivery):
        producer.produce(topic, key=key, value=value, on_delivery=on_delivery)
        producer.poll(0)

    def flush(self, producer):
        producer.flush(60)

    def make_consumer(self, args, group):
        from confluent_kafka import Consumer
        cfg = {
            "bootstrap.servers": args.bootstrap,
            "group.id": group,
            "auto.offset.reset": "earliest",
            "enable.auto.commit": False,
            "fetch.min.bytes": 1,
            "fetch.max.bytes": 50 * 1024 * 1024,
        }
        c = Consumer(cfg)
        c.subscribe([args.topic])
        return c

    def poll(self, consumer, timeout_ms):
        msg = consumer.poll(timeout_ms / 1000.0)
        if msg is None:
            return None
        if msg.error():
            return None
        return msg.value(), msg.timestamp()[1]


class KafkaPythonImpl(KafkaImpl):
    name = "kafka-python"

    def make_producer(self, args):
        from kafka import KafkaProducer
        return KafkaProducer(
            bootstrap_servers=args.bootstrap.split(","),
            acks=args.acks,
            linger_ms=args.linger,
            batch_size=args.batch,
            compression_type=None if args.compression == "none" else args.compression,
        )

    def produce(self, producer, topic, key, value, on_delivery):
        future = producer.send(topic, key=key, value=value)
        future.add_callback(lambda md, cb=on_delivery: cb(None, _Md(md)))
        future.add_errback(lambda exc, cb=on_delivery: cb(exc, None))

    def flush(self, producer):
        producer.flush(60)

    def make_consumer(self, args, group):
        from kafka import KafkaConsumer
        return KafkaConsumer(
            args.topic,
            bootstrap_servers=args.bootstrap.split(","),
            group_id=group,
            auto_offset_reset="earliest",
            enable_auto_commit=False,
            fetch_min_bytes=1,
            fetch_max_bytes=50 * 1024 * 1024,
        )

    def poll(self, consumer, timeout_ms):
        records = consumer.poll(timeout_ms=timeout_ms, max_records=500)
        msgs = []
        for tp, batch in records.items():
            for m in batch:
                msgs.append((m.value, m.timestamp))
        return msgs or None


class _Md:
    def __init__(self, md):
        self._md = md
    def topic(self): return self._md.topic
    def partition(self): return self._md.partition
    def offset(self): return self._md.offset


# ------------------------------- Producer 压测 -------------------------------

@dataclass
class ProducerStats:
    sent: int = 0
    failed: int = 0
    bytes_: int = 0
    latencies_us: List[int] = field(default_factory=list)
    start_ts: float = 0.0
    end_ts: float = 0.0


def run_producer(impl: KafkaImpl, args) -> ProducerStats:
    producer = impl.make_producer(args)
    stats = ProducerStats()
    sent_ts: dict[int, int] = {}
    lock = threading.Lock()

    def on_delivery(err, msg):
        with lock:
            now = time.perf_counter_ns()
            if err is not None:
                stats.failed += 1
                return
            seq_key = id(msg) if msg is not None else None
            t0 = sent_ts.pop(seq_key, None)
            if t0 is not None:
                stats.latencies_us.append((now - t0) // 1000)

    payload = (b"x" * args.size)
    print(f"\n[Producer] 开始发送 {args.records:,} 条消息,每条 {args.size} B …")
    stats.start_ts = time.perf_counter()

    for i in range(args.records):
        key = f"key-{i % max(1, args.keys)}".encode()
        sent_ts[id(payload) + i] = time.perf_counter_ns()
        try:
            impl.produce(producer, args.topic, key, payload,
                         lambda err, msg, k=id(payload) + i: _record(err, k, sent_ts, stats, lock))
            stats.sent += 1
            stats.bytes_ += args.size + len(key)
        except BufferError:
            time.sleep(0.001)
            continue

    impl.flush(producer)
    stats.end_ts = time.perf_counter()
    return stats


def _record(err, k, sent_ts, stats, lock):
    with lock:
        now = time.perf_counter_ns()
        t0 = sent_ts.pop(k, None)
        if err is not None:
            stats.failed += 1
            return
        if t0 is not None:
            stats.latencies_us.append((now - t0) // 1000)


def report_producer(stats: ProducerStats, args) -> None:
    dur = stats.end_ts - stats.start_ts
    if dur <= 0:
        print("⚠️  Producer 时长为 0")
        return
    msg_rate = stats.sent / dur
    mb_rate = stats.bytes_ / dur / 1024 / 1024
    lats = sorted(stats.latencies_us)
    print("=" * 72)
    print(" Producer 报告")
    print("-" * 72)
    print(f"  records sent : {stats.sent:,}  (failed {stats.failed})")
    print(f"  duration     : {dur:.2f} s")
    print(f"  msg rate     : {msg_rate:,.0f} msg/s")
    print(f"  byte rate    : {mb_rate:,.2f} MB/s")
    if lats:
        print(f"  latency avg  : {statistics.mean(lats)/1000:.2f} ms")
        print(f"  latency p50  : {lats[len(lats)//2]/1000:.2f} ms")
        print(f"  latency p95  : {lats[int(len(lats)*0.95)]/1000:.2f} ms")
        print(f"  latency p99  : {lats[int(len(lats)*0.99)]/1000:.2f} ms")
        print(f"  latency max  : {lats[-1]/1000:.2f} ms")
    print("=" * 72)


# ------------------------------- Consumer 压测 -------------------------------

@dataclass
class ConsumerStats:
    received: int = 0
    bytes_: int = 0
    start_ts: float = 0.0
    end_ts: float = 0.0


def run_consumer(impl: KafkaImpl, args, expected: int) -> ConsumerStats:
    group = f"bench-{uuid.uuid4().hex[:8]}"
    consumer = impl.make_consumer(args, group)
    stats = ConsumerStats()
    print(f"\n[Consumer] 开始消费,期望 {expected:,} 条 …")
    stats.start_ts = time.perf_counter()

    last_print = stats.start_ts
    while stats.received < expected:
        result = impl.poll(consumer, 1000)
        if not result:
            if time.perf_counter() - last_print > 5:
                print(f"  …已消费 {stats.received:,} / {expected:,}")
                last_print = time.perf_counter()
            if time.perf_counter() - stats.start_ts > args.consume_timeout:
                print("  ⚠️  消费超时,提前结束")
                break
            continue
        if isinstance(result, list):
            for value, _ts in result:
                stats.received += 1
                stats.bytes_ += len(value) if value else 0
        else:
            value, _ts = result
            stats.received += 1
            stats.bytes_ += len(value) if value else 0

    stats.end_ts = time.perf_counter()
    try:
        consumer.close()
    except Exception:
        pass
    return stats


def report_consumer(stats: ConsumerStats) -> None:
    dur = stats.end_ts - stats.start_ts
    if dur <= 0:
        print("⚠️  Consumer 时长为 0")
        return
    print("=" * 72)
    print(" Consumer 报告")
    print("-" * 72)
    print(f"  records read : {stats.received:,}")
    print(f"  duration     : {dur:.2f} s")
    print(f"  msg rate     : {stats.received/dur:,.0f} msg/s")
    print(f"  byte rate    : {stats.bytes_/dur/1024/1024:,.2f} MB/s")
    print("=" * 72)


# --------------------------------- main ---------------------------------

def main() -> None:
    parser = argparse.ArgumentParser(description="Kafka Producer/Consumer 端到端吞吐压测")
    parser.add_argument("--bootstrap", default=os.getenv("KAFKA_BOOTSTRAP", "localhost:9092"))
    parser.add_argument("--topic", default="learn.bench")
    parser.add_argument("--records", type=int, default=100_000)
    parser.add_argument("--size", type=int, default=1024, help="单条消息字节数")
    parser.add_argument("--keys", type=int, default=1000, help="key 基数(控制分区分布)")
    parser.add_argument("--batch", type=int, default=16384, help="batch.size 字节")
    parser.add_argument("--linger", type=int, default=0, help="linger.ms")
    parser.add_argument("--compression", default="none",
                        choices=["none", "gzip", "snappy", "lz4", "zstd"])
    parser.add_argument("--acks", default="1", choices=["0", "1", "all"])
    parser.add_argument("--no-consume", action="store_true", help="只压 Producer")
    parser.add_argument("--consume-timeout", type=int, default=120)
    args = parser.parse_args()

    impl = KafkaImpl.detect()
    print(f"使用客户端实现:{impl.name}")
    print(f"目标 Topic    :{args.topic} @ {args.bootstrap}")
    print(f"配置:batch={args.batch} linger={args.linger}ms "
          f"compression={args.compression} acks={args.acks}")

    p_stats = run_producer(impl, args)
    report_producer(p_stats, args)

    if not args.no_consume:
        c_stats = run_consumer(impl, args, expected=p_stats.sent)
        report_consumer(c_stats)

    # 给一些直觉化建议
    print("\n💡 调参建议:")
    if args.linger == 0 and args.batch <= 16384:
        print("  - 你目前是『极致低延迟』模式,吞吐受限。如果是日志采集 / 数仓场景,把 batch 提到 64KB+,linger 提到 20ms+。")
    if args.compression == "none" and args.size > 256:
        print("  - 没开压缩。文本类消息开启 lz4 / zstd 通常吞吐翻倍,CPU 涨 5~15%。")
    if args.acks == "0":
        print("  - acks=0 数据可能丢,仅适用于「丢一点没关系」的指标埋点。")
    if args.acks == "all":
        print("  - acks=all 是金标准。配合 RF=3 + min.insync.replicas=2 才算真的不丢。")


if __name__ == "__main__":
    sys.exit(main())

sequential_vs_random.py ↗ · throughput_bench.py ↗