主题
第 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)[典型值] | 5× |
| NVMe SSD(企业级) | ~6000 MB/s [典型值] | ~1500 MB/s(约 350K IOPS)[典型值] | 4× |
数据来源: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
Disksendfile 的工作流:
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 + write | 2 | 2 | 4 | 2 | 2 |
sendfile(pre 2.4) | 2 | 1 | 3 | 1 | 1 |
sendfile + DMA gather(2.4+) | 2 | 0 | 2 | 1 | 1 |
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/write | mmap |
|---|---|---|
| 数据流 | 用户 buffer 与 PageCache 互相 copy | PageCache(不复制到用户态) |
| 系统调用次数 | 每次 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 / zstd5 种压缩对比 [典型值]:
| 算法 | 压缩率 | 压缩速度 | 解压速度 | CPU 占用 | 推荐场景 |
|---|---|---|---|---|---|
none | 1.0x | - | - | 0 | 局域网内 + 已经是二进制 |
gzip | 5~7x | 慢 | 中 | 高 | 老遗留集群 |
snappy | 2~3x | 快 | 快 | 低 | Kafka 默认很久 |
lz4 | 3~4x | 极快 | 极快 | 低 | 生产首选 |
zstd | 4~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.ms | batch 平均大小 | 吞吐 (MB/s) | p99 延迟 (ms) |
|---|---|---|---|
| 0 | 1 KB | 80 | 5 |
| 5 | 12 KB | 200 | 8 |
| 20 | 32 KB | 380 | 25 |
| 50 | 60 KB | 480 | 55 |
| 100 | 65 KB | 510 | 105 |
| 500 | 65 KB(达到 batch.size 上限) | 520 | 505 |
观察:
- linger 从 0 -> 20,吞吐提升 5 倍,延迟仅增加 20ms。
- linger 从 20 -> 100,吞吐只再提升 30%,延迟却涨 5 倍。
- 甜点(sweet spot)通常在 5 ~ 30ms,业务能接受的延迟内取最大。
7.3 业务取舍
| 业务场景 | 推荐 linger.ms |
|---|---|
| 实时风控(要求 <= 100ms 端到端) | 5 ~ 10 |
| 一般日志聚合 | 20 ~ 50 |
| 离线批量 ETL(容忍秒级) | 100 ~ 500 |
| 低吞吐高优先级(订单事件) | 0 ~ 5 |
8. 与 RabbitMQ / RocketMQ / Pulsar 对比
| 维度 | Kafka | RabbitMQ | RocketMQ | Pulsar |
|---|---|---|---|---|
| 主要存储模式 | 顺序写 + PageCache + sendfile | Erlang Mnesia + 内存为主 | CommitLog 顺序写 + ConsumeQueue 索引 | Bookkeeper Ledger 多副本写 |
| 零拷贝 | sendfile | 否 | sendfile | sendfile |
| Broker 是否在堆里缓存消息 | 否(依赖 PageCache) | 是(内存 + 持久化) | 否(PageCache) | 否 |
| 批量 | Producer 端攒 batch | 单条 ack 多 | 单条发送、batch 较少 | Producer 攒 batch |
| 压缩粒度 | 整个 batch | 单条 | 单条 | 整个 entry |
| 典型单 Broker 吞吐 | 几十万 msg/s, 几百 MB/s | 几万 msg/s | 几十万 msg/s | 数十万 msg/s |
| 延迟 p99 | 5ms ~ 50ms(取决 linger.ms) | 1ms ~ 10ms | 5ms ~ 30ms | 5ms ~ 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 为什么这么快?
考察点:综合理解高吞吐底层原理。
答案:
- 顺序写:消息只 append 到 active segment 的
.log末尾。HDD 顺序写比随机写快 100x,SSD 上也快 5x。 - PageCache:Broker 不在 JVM 堆里缓存消息,靠 OS 的 PageCache 缓存。优势:
- 避免 GC 暂停。
- Broker 重启不丢缓存(PageCache 在内核里)。
- 同一份数据不在堆和 PageCache 里同时存两份。
- 零拷贝:Consumer 拉数据走
sendfile(在 Kafka 里通过FileChannel#transferTo)。从read+write的 4 次拷贝 + 2 次上下文切换,降到 2 次 DMA 拷贝 + 1 次上下文切换。 - mmap:
.index和.timeindex用 mmap,二分查找是纯内存指针运算。 - 批量 + 压缩:Producer 攒 batch、整 batch 压缩,Broker 原样落盘和转发,Consumer 端才解压。压缩算法首选 lz4(吞吐最优,CPU 低)。
- 加分项:
- SSL 不能 sendfile,会让吞吐腰斩。
- Producer 端
linger.ms=20通常是吞吐/延迟最优拐点。 - 给 PageCache 预留足够内存(建议 JVM heap <= 8GB)。
Q2:sendfile 怎么做到「零拷贝」?比传统 read+write 少了几次拷贝?
考察点:操作系统、Linux 系统调用。
答案:
- 传统 read + write 一共 4 次拷贝、2 次上下文切换:
- 磁盘 -> PageCache(DMA)
- PageCache -> 用户 buffer(CPU)
- 用户 buffer -> socket buffer(CPU)
- socket buffer -> NIC(DMA)
- sendfile(pre 2.4 内核)3 次拷贝、1 次上下文切换:
- 磁盘 -> PageCache(DMA)
- PageCache -> socket buffer(CPU)
- socket buffer -> NIC(DMA)
- sendfile + DMA gather(2.4+ 内核)2 次 DMA 拷贝、1 次上下文切换:
- 磁盘 -> PageCache(DMA)
- PageCache -> NIC(DMA gather,NIC 直接从 PageCache 读)
- 「零拷贝」中的「零」是指零 CPU 拷贝,不是物理上没有拷贝。
- Kafka 里的应用:
- Consumer fetch 走
FileChannel#transferTo-> JVM -> 底层 sendfile。 - Producer 写 Broker 不走 sendfile(数据来自网络 buffer 不是文件)。
- SSL 连接走不了 sendfile(数据要先到用户态加密)。
- Consumer fetch 走
- 加分项:提到 Java 21 之前
transferTo在 SSL 下会 fallback 到普通 read/write,21 之后引入 io_uring 有进一步优化。
Q3:Kafka Broker 为什么不在 JVM 堆里缓存消息?
考察点:PageCache、JVM、设计哲学。
答案:
- JVM 对象开销大:
byte[1024]实际占用约 1052 字节,包装成 Message 对象后接近 1.5KB。 - GC 是定时炸弹:堆里几十 GB 的对象,每次 Major GC 暂停几秒到几十秒,Kafka 的吞吐目标受不了。
- 数据双份浪费:同一份消息如果堆里和 PageCache 里都存,纯属浪费。OS 在内存压力下还得去淘汰一份,反而恶化。
- Broker 重启不丢缓存:JVM 堆缓存进程退出就没;PageCache 在内核里,进程重启后还在。Kafka 滚动升级时几乎无吞吐损失。
- OS 已经把缓存做得很好:PageCache 自动 LRU、自动预读、自动写回,Kafka 没必要重新发明。
- 配置佐证:官方建议 Heap 6~8GB、剩下都让出给 PageCache;
log.flush.interval.messages/log.flush.interval.ms默认都是无穷大,意思是「不用主动 flush,让 OS 决定」。 - 加分项:提到 Kafka 是 Scala 写的(JVM 体系)但消息从来不进 Scala 对象层,全程 ByteBuffer,本质就是「让 OS 帮我做一切」。
Q4:Kafka 顺序写真的比随机写快两个数量级吗?给真实数字。
考察点:对存储硬件的认知。
答案:
- HDD(机械盘 7200 rpm SATA)[典型值]:
- 顺序写:~150 MB/s
- 随机写 4KB:~1 MB/s(约 100 IOPS)
- 倍数:~150x,因为机械臂寻道(5ms) + 旋转延迟(4ms) 占绝对主导。
- SATA SSD [典型值]:
- 顺序:~500 MB/s
- 随机 4KB:~50 MB/s
- 倍数:~10x。SSD 没有寻道,但 NAND 的 GC 写放大让随机写慢一截。
- NVMe SSD(消费级)[典型值]:
- 顺序:~3 GB/s
- 随机 4KB:~600 MB/s
- 倍数:~5x。NVMe 仍然偏好顺序,但差距收窄。
- Kafka 怎么用:永远 append 到
.log末尾,永远不修改已有消息。删除按 segment 整体 unlink,不是按消息删。 - 加分项:提到 RDBMS 的 redo log 也是顺序写,但数据页随机写让总体仍偏慢;Kafka 没有「数据页」,写完就完事。
Q5:mmap 在 Kafka 哪里用?为什么 .log 不 mmap?
考察点:mmap 原理、Kafka 存储结构。
答案:
- 应用位置:仅
.index和.timeindex文件用 mmap(org.apache.kafka.common.utils.OffsetIndex)。 - mmap 优势:
- 把文件直接映射到进程虚拟地址空间,访问像访问数组指针。
- 二分查找时没有 read 系统调用,只有第一次 page fault。
- 写新索引项也是直接写 mmap,不需要 write 系统调用。
.log为什么不 mmap:- 单文件 1GB,mmap 进程地址空间不灵活(32 位下根本装不下,64 位也会让 PageCache 锁住一大片)。
.log的访问模式是「顺序读写」,read/write + sendfile 已经很高效,mmap 收益不大。- mmap 写盘是异步的,崩溃可能丢最后几个 page,但
.index可以从.log重建,.log不能丢。
- 崩溃恢复:
.index丢了几条索引没关系,Broker 启动时会用.log重新构建索引。 - 加分项:提到 Java NIO 的
FileChannel.map()是 mmap 入口;MappedByteBuffer在 GC 后才能真正 unmap,所以 Broker 关闭时.index文件可能要等几秒才能被 OS 真正释放。
Q6:批量 + 压缩在 Producer / Broker / Consumer 三端怎么协同?
考察点:Kafka 数据流。
答案:
- Producer 端:
- 多条消息进入 RecordAccumulator 攒 batch(按 partition 维度)。
- 凑够
batch.size(默认 16KB)或linger.ms到时间,整 batch 用 lz4/zstd 压缩。 - 一次网络 send 把整个压缩 batch 发给 Broker。
- Broker 端:
- 收到压缩 batch 后完全不解压,直接 append 到
.log末尾。 - 既省 Broker CPU,又省磁盘空间。
- Topic 端建议设
compression.type=producer透传 Producer 压缩。
- 收到压缩 batch 后完全不解压,直接 append 到
- Follower 副本:拉到的也是压缩 batch,同样不解压直接 append。
- Consumer 端:
- 通过
sendfile拿到压缩 batch(连压缩字节都不变)。 - 在 Consumer 进程解压、拆 batch 成单条 record 交给业务。
- 通过
- 化学反应:一个 batch 内消息往往字段名、格式重复,压缩率比单条压缩高得多(5~10x)。
- 配置推荐:
compression.type=lz4(吞吐最优、CPU 低)batch.size=65536(64KB)linger.ms=20
- 加分项:提到 Broker 如果设
compression.type=lz4且 Producer 是 gzip,Broker 会重新压缩(解压再压),吞吐和 CPU 都受损;所以默认producer模式最好。
Q7:linger.ms 调高会怎样?给一个实测对比。
考察点:吞吐 vs 延迟取舍。
答案:
- 作用:Producer 在没凑够
batch.size时,最多等linger.ms毫秒再发,期望期间能凑出更大的 batch。 - 公式:
- 吞吐 ~ batch 平均大小(更大 batch -> 每次系统调用 + 压缩 + 网络携带的消息更多)。
- 延迟 ~ linger.ms + 网络 RTT + Broker 处理(5~10ms)。
- 实测对比 [典型值](同一集群、同一 producer):
linger=0, batch.size=16K-> 吞吐 80MB/s, p99 延迟 5mslinger=5-> 200MB/s, 8mslinger=20-> 380MB/s, 25mslinger=50-> 480MB/s, 55mslinger=100-> 510MB/s, 105ms
- 拐点:通常在 5~30ms。再大吞吐增长缓慢,延迟却线性涨。
- 业务推荐:
- 实时风控/订单 ->
linger=5以下。 - 日志聚合 ->
linger=20~50。 - 离线 ETL ->
linger=100~500。
- 实时风控/订单 ->
- 加分项:提到
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:用 Pythonos.write模拟顺序写 vs 随机写文件吞吐对比。08_performance/code/throughput_bench.py:完整 Producer -> Consumer 端到端压测脚本。
延伸阅读
- 第 7 章 存储与日志格式 —— 「顺序写」「mmap」依赖的存储模型。
- 第 4 章 Producer 深入 ——
linger.ms、batch.size、compression.type的更细致讨论。 - 第 5 章 Consumer 深入 ——
fetch.min.bytes、fetch.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())