Kafka 为什么能实现高吞吐?顺序写、批处理、页缓存与零拷贝
Kafka 的高吞吐并不能简单归因于“顺序写磁盘”。真正的答案是一套相互配合的设计:分区日志降低随机访问,生产者批量发送并压缩,Broker 借助页缓存与零拷贝传输,消费者主动拉取并批量处理,最后通过增加分区横向扩展。
一、分区日志:把问题变成追加写
一个 Topic 被拆成多个 Partition,每个分区是一条有序、不可变的追加日志。新消息只写到尾部,并获得递增 offset:
Partition 0: [0][1][2][3][4] -> append
Partition 1: [0][1][2] -> append
磁盘顺序 I/O 的吞吐远高于大量随机寻址。Kafka 还把大日志切成多个 segment,每个 segment 配有索引,查找时先定位 segment,再通过稀疏索引逼近目标 offset。
二、批处理减少系统调用
生产者不是每条消息都发一次请求,而是先按分区聚合为 RecordBatch:
Properties p = new Properties();
p.put("bootstrap.servers", "kafka:9092");
p.put("batch.size", 64 * 1024);
p.put("linger.ms", 10);
p.put("compression.type", "lz4");
batch.size 控制批次目标大小,linger.ms 允许生产者短暂等待更多消息。批量能摊薄网络往返、协议头、校验和系统调用成本,但会增加少量排队延迟。
版本默认值不能凭旧经验记忆:Kafka 4.0 把生产者 linger.ms 默认值从 0ms 调整为 5ms,以换取更高效批处理。升级 3.x 到 4.x 时应对 P99 延迟、批次大小和吞吐重新压测,而不是机械覆盖成旧默认值。
调优不能孤立地把批次调大。低流量分区可能始终凑不满批次,过大的缓冲还会增加客户端内存压力。
三、批量压缩为什么反而更快
Kafka 对整个批次压缩,相似消息在一起时压缩率通常优于逐条压缩。虽然压缩消耗 CPU,却减少网络与磁盘字节数,整体吞吐常常更高。
常见选择:
| 算法 | 特点 | 场景 |
|---|---|---|
| lz4 | 速度与压缩率均衡 | 通用低延迟 |
| snappy | 快、成熟 | 通用场景 |
| zstd | 压缩率更高 | 网络或存储更敏感 |
| gzip | 压缩率高但 CPU 较重 | 不强调低延迟 |
需要让生产者、Broker 与消费者版本都支持所选算法,并用真实消息体压测。
四、页缓存:Kafka 为什么不自己维护复杂缓存
写入日志时,数据通常先进入操作系统 page cache,再由内核异步刷盘。读取热点数据也大概率直接命中页缓存,无需真正访问磁盘。
Kafka 倾向把内存交给操作系统管理,而不是在 JVM 堆中维护大块消息缓存,带来几个好处:
- 避免巨大 Java 堆导致的 GC 压力。
- Broker 重启后,操作系统页缓存仍可能保留热点。
- 内核能统一协调读写与回收策略。
但页缓存不等于数据已经持久化。可靠性还取决于副本、确认级别和 ISR,而不是“写入接口返回”这一件事。
五、零拷贝减少用户态搬运
传统文件发送大致需要:磁盘 → 内核页缓存 → 用户态缓冲区 → Socket 缓冲区 → 网卡。Kafka 在合适路径下利用 sendfile 等机制,让数据从页缓存直接进入网络发送路径,减少用户态复制和上下文切换。
零拷贝并不是完全没有任何硬件复制,而是避免不必要的内核态与用户态往返。若消息需要在 Broker 端解压、修改、重新编码,或 TLS/运行时组合无法使用同等内核直通路径,就不能假设仍有相同收益,应该以实际 CPU profile 和网络指标验证。
六、消费者拉取与批量消费
Kafka 采用 pull 模型,消费者按照自身能力请求数据:
@KafkaListener(topics = "orders", batch = "true")
public void consume(List<ConsumerRecord<String, OrderEvent>> records) {
orderHandler.handleBatch(records);
}
消费者可通过 fetch.min.bytes、fetch.max.wait.ms 等参数平衡吞吐与延迟。批量拉取减少请求次数,也更容易批量写数据库。
不过批量越大,单次处理时间越长。若超过 max.poll.interval.ms,消费者可能被踢出组并触发 rebalance,反而导致重复处理和吞吐抖动。
七、分区带来横向扩展
一个分区在同一个消费者组内只能由一个消费者实例消费,但不同分区可以并行。因此吞吐上限与分区数、Broker 分布和消费者并行度密切相关。
分区也有成本:更多文件句柄、更多复制流量、更长的故障恢复与选主时间。不能为了“以后可能扩容”无上限增加分区,而且 Kafka 通常只支持增加分区,增加后基于 key 的映射可能变化。
八、副本与吞吐的权衡
生产者使用 acks=all 时,Leader 要等待 ISR 中满足条件的副本确认。可靠性更强,但跨 Broker 复制会增加延迟与网络开销。
典型生产组合:
acks=all
enable.idempotence=true
min.insync.replicas=2
replication.factor=3
不能只在客户端设置 acks=all,还要确保 Topic 副本数和 Broker 的 min.insync.replicas 匹配。min.insync.replicas=2 表示 ISR 少于 2 时拒绝写入;acks=all 则等待当时 ISR 中的副本确认,并不只是固定等待两个副本。可靠性与可用性还要结合副本布局、机架感知和故障域评估。
Kafka 4.0 已移除 ZooKeeper 模式,仅支持 KRaft。KRaft 改变元数据与控制面,不会推翻本文的日志追加、批处理和页缓存这些数据面原理;但升级时仍应重新验证选主恢复、消费者再均衡和运维工具链。
九、常见性能误区
- 只提高分区数,不检查消费者和 Broker 是否成为瓶颈。
- 把 JVM 堆调得极大,挤占操作系统页缓存。
- 每条消息同步
get()等待发送结果,破坏批处理。 - 消费者逐条写数据库,吞吐被数据库往返限制。
- 只看平均延迟,不看 P99、rebalance 和积压增长率。
- 压测消息很小且内容重复,得到不真实的压缩收益。
十、性能排查清单
- 生产端批次大小、压缩率和请求延迟是多少?
- 分区是否均衡,是否存在热点 key?
- Broker 磁盘吞吐、网络、页缓存命中与 CPU 是否饱和?
- ISR 是否频繁收缩,副本复制是否落后?
- 消费端 poll 与处理耗时是否接近超时阈值?
- 消费者实例数是否超过有效分区数?
- 积压是持续增长还是流量峰值后的正常回落?
参考基线
总结
Kafka 的高吞吐来自端到端批量化与操作系统友好设计:生产者聚合并压缩,Broker 顺序追加、使用页缓存和零拷贝,消费者批量拉取,分区再提供并行扩展。每项机制都在吞吐、延迟、可靠性和资源占用之间做权衡,调优必须从完整链路和真实业务负载出发。