Skip to content

Kafka 为什么快:吞吐来自批处理和分区,而不是某一个技巧 ​

「顺序写、页缓存、零拷贝」是常见的回答,但把它们罗列一遍,解释不了为什么同一个集群换一组生产者参数,吞吐能差 40 倍;也解释不了为什么加了消费者,消费速度却一点没变。Kafka 的吞吐来自一组相互配合的设计,每一项都有它的前提和代价。

同一个三节点集群、同样 20 万条 256 字节的消息,关闭批处理时生产者每秒只能发出约 9450 条,平均延迟好几秒;默认参数下约 18 万条;把批再调大,约 37 万条。消费端也一样:6 个消费者消费 1 个分区的 topic,只有 1 个在干活;换成 6 个分区,耗时从 7.1 秒降到 1.2 秒;可如果 80% 的消息用了同一个 key,又回到 5.9 秒。

本文在 Kafka 4.3.1 的三节点 KRaft 集群上,逐一测量批处理、压缩、确认级别、分区并行和热点 key 对吞吐与延迟的影响。

一、先说结论 ​

  • 批处理是生产端吞吐的主要来源:一个请求带 1 条和带约 1000 条消息,吞吐差约 40 倍。linger.ms 让生产者稍等一会儿凑批,用一点延迟换大幅吞吐。
  • 追加写与页缓存降低了 broker 的单条成本:分区是只追加的日志,读写都是顺序的,热数据留在操作系统的页缓存里,不需要 JVM 堆再缓存一份。
  • 压缩省的是磁盘和网络:同样 20 万条 JSON,磁盘上的日志从 31 MB 降到 lz4 的 7.5 MB、zstd 的 4.2 MB;按批压缩,批越大压缩率越高。
  • acks 是用延迟换可靠:3 副本下 acks=all 的平均延迟约是 acks=1 的 3.6 倍,p99 从 33ms 到 183ms。只靠调低 acks 得到的吞吐,不是同一可靠性下的性能提升。
  • 分区决定并行度,key 决定分区:1 个分区只能被一个消费者处理;80% 的消息用同一个 key 时,一个分区分到 83% 的消息,加再多消费者也没用。

二、一条消息在 Kafka 里经过什么 ​

生产者按分区攒批 · 整批压缩batch.size · linger.msbroker(分区 leader)校验后追加到日志末尾页缓存 · 不重新压缩消费者按 offset 顺序整批读取自己记录消费位置一批一批副本 follower × 2按批从 leader 拉取acks=all 要等它们关闭批处理 约 9,450 条/秒 · 默认 约 18.2 万 · batch.size=256KB、linger.ms=20 约 36.6 万
图 1 · 生产者按分区攒批并整批压缩;broker 校验后把整批追加到日志末尾,热数据留在页缓存;消费者按 offset 顺序整批读取。批是贯穿整条链路的单位,实测一个请求带约 1000 条时吞吐约为带 1 条时的 40 倍
  • 生产者把发往同一分区的消息攒成批(record batch),整批压缩、整批发送。
  • broker 校验整批数据后追加到分区日志的末尾,不重新压缩、不拆开改写,副本从 leader 按批拉取。
  • 消费者按 offset 顺序读取整批数据,自己记录消费到哪里。

这条链路上,「一批」是贯穿始终的单位:一次网络请求、一次校验、一次磁盘写入、一次压缩,承载的是几十到上千条消息。单条消息越小,批处理的收益越明显。

三、批处理:一个请求带多少条 ​

同一个 topic(1 分区、1 副本)、acks=1,三组参数各跑 3 轮取中位数:

参数吞吐每个请求平均条数批大小平均
batch.size=0(关闭批处理)约 9,450 条/秒1.0326 字节
默认(batch.size=16384,linger.ms=5)约 18.2 万条/秒61.016,224 字节
batch.size=262144,linger.ms=20约 36.6 万条/秒980.4260,781 字节

关闭批处理后,平均延迟反而达到数秒:每条消息一个请求,发送跟不上生产,消息在客户端缓冲区里排起了长队。所以「批处理增加延迟」只在低流量时成立;流量一高,不批处理的延迟更差。

  • batch.size 是每个分区一个批的上限,满了立即发送;
  • linger.ms 是批没满时最多等多久,Kafka 4.0 起默认值从 0 改成了 5ms;
  • 批大,单条消息的平均开销低,但内存占用和尾部延迟会上升,需要在自己的消息大小和延迟目标下压测。

四、追加写、页缓存与零拷贝 ​

这三项是 broker 端单条成本低的原因,本文没有单独测量,结论来自 Kafka 官方文档的设计说明:

  • 追加写:每个分区是一组只追加的日志段文件,写入总是写在末尾,读取按 offset 顺序进行。顺序 I/O 在机械盘和 SSD 上都比随机 I/O 便宜得多,也让预读和合并写入更有效。
  • 页缓存:broker 不在 JVM 堆里缓存消息,而是依赖操作系统的页缓存。刚写入的数据通常还在内存里,消费者追着读时基本不访问磁盘;broker 重启后缓存也还在,而且不增加 GC 压力。代价是需要给操作系统留足内存,其他进程挤占页缓存时,消费延迟会明显上升。
  • 零拷贝:没有 TLS 时,broker 可以用 sendfile 把日志文件的数据直接从页缓存发到网卡,不经过用户态复制。开启 TLS 后数据必须在用户态加密,这条路径就不再适用。所以「零拷贝」的收益取决于是否加密,不能笼统地算作 Kafka 快的首要原因。

broker 能做到不重新压缩,前提是生产者和 broker 的消息格式、压缩配置一致;topic 设置了和生产者不同的 compression.type 时,broker 要解压再重新压缩,额外消耗 CPU。

五、压缩:省磁盘和带宽 ​

消息体是字段名重复、取值有限的业务 JSON,三个 topic 各写 20 万条:

压缩方式磁盘上的日志
none31.0 MB
lz47.5 MB
zstd4.2 MB

压缩按批进行,批越大、消息之间重复的内容越多,压缩率越高,这是批处理的又一个收益。压缩省下的是磁盘空间、副本复制的带宽和消费者拉取的带宽,代价是生产者和消费者的 CPU。

实验中三种方式的生产吞吐只跑了一轮,波动很大(另一次运行里 lz4 反而比不压缩更快),本文不据此比较速度。选择压缩算法时,用自己的消息体和机器测一次:通常 lz4 的 CPU 开销低,zstd 的压缩率高。

六、acks:确认到哪一步 ​

3 副本、min.insync.replicas=2 的 topic:

acks固定 2 万条/秒时的平均延迟p99不限速的吞吐
02.1ms12ms约 18.7 万条/秒
13.4ms33ms约 17.3 万条/秒
all12.4ms183ms约 7.8 万条/秒

acks=all 要等同步副本都写入才确认,多了一次副本拉取的往返,延迟和尾部延迟都明显上升。三种设置代表三种不同的可靠性:acks=0 发出去就不管,acks=1 在 leader 写入后、副本同步前 leader 宕机会丢,acks=all 配合 min.insync.replicas 才能在一个 broker 故障时不丢已确认的消息。各个环节分别由什么保证,见 Kafka 不丢、不重与 Exactly Once。

评估性能时,吞吐和可靠性要一起写:「acks=all 下 7.8 万条/秒」和「acks=0 下 18.7 万条/秒」不是同一个系统的两个优化结果。

七、分区:并行度的上限 ​

一个分区在同一个消费组内同一时刻只分给一个消费者,所以分区数是消费并行度的上限。3000 条消息、每条处理 2ms、6 个消费者:

text
1 个分区:从第一条到最后一条 7.1 秒,真正处理过消息的消费者 1 个
6 个分区:从第一条到最后一条 1.2 秒,真正处理过消息的消费者 6 个

只有分区够多,加消费者才有意义。反过来,分区也不是越多越好:每个分区都有自己的日志文件、索引、副本同步和 leader 选举,分区过多会增加 broker 的元数据、文件句柄和故障恢复时间,重平衡也更慢,见 Kafka 消费组与重平衡。分区数从目标吞吐和单分区的实测吞吐估算,再留出增长余量。

7.1 热点 key 把并行度打回一个分区 ​

有 key 的消息按 key 的哈希决定分区,同一个 key 总在同一个分区,这保证了同一个 key 的消息有序。代价是 key 分布不均时,分区也不均:

1 个分区6 个分区,key 均匀6 个分区,80% 同一个 keyP0 · 3000 条~1000~1000~1000~1000~1000~1000~200~200~200~2004976~2007.1 秒 · 1 个消费者在处理1.2 秒 · 6 个消费者5.9 秒 · 由热点分区决定每条处理 1—2ms;同一个 key 总在同一个分区,保证了它的顺序,也集中了它的负载
图 2 · 6 个消费者:1 个分区时只有 1 个在处理;6 个分区且 key 均匀时 6 个一起处理;80% 的消息用同一个 key 时,一个分区分到 83%,总耗时由它决定
text
key 均匀分布:6000 条分到 6 个分区 [1026, 1035, 999, 989, 970, 981],最大分区占 17%;6 个消费者:1.2 秒
80% 的消息 key 都是 event-42:6 个分区 [202, 204, 214, 196, 4976, 208],最大分区占 83%;6 个消费者:5.9 秒

热点分区上的 4976 条只能由一个消费者处理,总耗时由它决定。常见的处理方式:

  • 重新选择 key:只有需要保证顺序的实体才用作 key,比如用订单号而不是商户号;
  • 给热点 key 加后缀分散到多个分区,代价是这个 key 内部不再全局有序;
  • 在消费者内部,把同一分区的消息按更细的 key 分到多个线程处理。

八、常见瓶颈与排查 ​

现象可能的原因看什么
生产吞吐上不去,CPU 不高批太小,请求太多生产者指标 records-per-request-avg、batch-size-avg
生产者延迟很高发送缓冲区排队、acks=all 等待副本record-queue-time-avg、request-latency-avg
消费 lag 持续增长,broker 很空闲消费处理慢,或分区太少每个分区的 lag、消费者处理耗时
某个分区 lag 特别大热点 key各分区的写入量
消费延迟突然上升页缓存被挤占,开始读磁盘broker 所在机器的内存与磁盘读

九、常见误区 ​

  • 「Kafka 快是因为零拷贝」:零拷贝只在不加密时成立,批处理和顺序 I/O 的影响更大。
  • 「批处理会增加延迟」:低流量时是;流量一高,不批处理的延迟反而更差,实测达到数秒。
  • 「调低 acks 就能提升性能」:换来的是可靠性下降,不是同一条件下的优化。
  • 「消费慢就加消费者」:消费者数超过分区数的部分是空闲的;热点 key 时加分区也没用。
  • 「分区越多越好」:分区带来元数据、文件句柄、复制和恢复时间的成本。

小结 ​

Kafka 的吞吐来自一组相互配合的设计:生产者把消息攒成批并按批压缩,broker 按批顺序追加、依靠页缓存服务读取,分区让生产和消费可以横向并行。每一项都有前提:批处理需要足够的流量或一点等待,压缩需要消息之间有重复,分区并行需要 key 分布均匀。调参数之前先看生产者指标和各分区的 lag,找到是哪一环限制了吞吐。


配套实验

参考资料

文章以 CC BY-NC-SA 4.0 授权 · 代码片段以 MIT 授权