系列目录
- MySQL 索引与慢查询:B+ 树如何减少扫描
- 声明式事务之下:InnoDB 的 MVCC 与锁
- 接入层:Nginx 反向代理与 OpenResty 的边界
- Web 容器与 Netty:线程模型之下的 IO 模型
- Redis(上):缓存用法与单线程模型的限制
- Redis(下):超出缓存用途的用法:锁、队列与排行榜
- 数据库访问层:连接池与 MyBatis 的显式 SQL
- Kafka(上):吞吐的来源是顺序 IO(本篇)
Kafka 处理峰值流量的原因
之前(Kafka 异步化与削峰:热点上报和最终一致性)分析了云盘用 Kafka 做异步化与削峰的两个链路。那篇落在业务接入层:本地事件表保证业务数据与待发消息原子写入、消费端用去重表做幂等、定时对账兜底。当时 Kafka 被当作可信基础设施接入,配置照搬公司 mafka 平台的默认值,按经验调生产者批次和消费者并发。当时的判断是,Kafka 能处理峰值流量。
当时没有弄清楚其中原因。相同请求量直接写入 MySQL 时,MySQL 会先达到容量上限;经 Kafka 削峰后,链路可以保持稳定。两者都落盘,Kafka 还多了一层网络投递,差异来自二者的存储和 IO 模式。2018 年补充 Kafka 设计资料时,这个问题才有了答案。
2016 第 6 篇介绍业务接入方式,包括事件表、幂等和对账;本篇说明 broker 的吞吐来源,不重复前文内容。
顺序写:append-only 日志把写入变成追加
磁盘的机械结构决定了顺序写远快于随机写。机械盘的一次写入要先寻道(磁头移动到目标磁道)再等旋转(目标扇区转到磁头下),这两段时间在毫秒量级;顺序写磁头不动,写入只受磁盘带宽限制。SSD 没有机械寻道,但随机写仍会触发更多的擦除放大和垃圾回收,顺序写在 SSD 上同样更友好。量级差:机械盘顺序写吞吐可达上百 MB/s,随机写常在个位数 MB/s。
Kafka 的存储模型是 append-only 日志。每个 partition 在 broker 上是一组日志段(log segment)文件:.log 存消息、.index 存 offset 到文件位置的稀疏索引、.timeindex 存时间戳到 offset 的索引。写入永远追加到当前活跃段的 .log 末尾,不修改已写入的消息。一段写满(log.segment.bytes 或 log.roll.hours 触发)就切到下一段。消费者按 offset 顺序读取,.index 帮它把 offset 换算成文件位置,再顺序读。
追加写同时影响写入和读取。写入侧没有更新、没有原地删除(删除靠 retention 策略整段丢弃或 log cleanup 的 compact),磁盘访问模式是纯顺序。读取侧消费者也是顺序读,对 page cache 友好。存储模型本身保证了访问模式是顺序的。
Kafka 与数据库的写入成本不同。数据库除了顺序写 redo log,还要在 buffer pool 里更新数据页、维护 B+ 树索引(索引页的写入随数据分布散布,带随机 IO)、写 undo log 支持回滚与 MVCC,再在提交时保证事务的持久性与隔离性。Kafka 的写入只是往日志末尾追加一条记录,没有索引页更新、没有事务协调。Kafka 是纯顺序追加;数据库同时承担顺序日志、随机索引维护和事务处理,因此二者的磁盘 IO 模式和成本不同。
Kafka 如何使用 page cache 缓存消息
Kafka 不在 JVM 堆里维护消息数据缓存。写入路径是:producer 把消息发给 broker,broker 把数据写入 page cache,由操作系统异步刷盘(log.flush.interval.messages / log.flush.interval.ms 可配,默认交给 OS)。读取路径是:consumer 拉取时,broker 从 page cache 命中数据直接返回。
不在堆里维护消息缓存有三个作用。第一,避免双份内存:消息数据已经在 page cache 里,再在 JVM 堆里存一份等于同一份数据占两份内存。第二,减少消息缓存带来的 GC 压力:堆内大缓存是大对象,GC 扫描和搬运的代价随堆增大上升;放在 page cache 里不进堆。第三,broker 重启后 page cache 仍在(只要 OS 没被重启),热数据不需要重新加载。
这对 JVM 参数配置的含义是:Kafka broker 的 -Xmx 不宜给太大,通常 6–8 GB 够用,机器剩余内存留给 page cache。给堆分太多内存会增加 GC 停顿时间,也会挤压 page cache 的可用空间。许多 Java 应用倾向于将更多内存分配给堆;Kafka broker 则需要限制堆大小,为 page cache 留出空间。
sendfile:数据不进用户态
2016 系列第 9 篇分析过文件到 socket 的数据路径(流对流的失败率、NIO 与零拷贝)。那篇在业务代码层面说明传统 read+write 的四次数据移动(DMA 磁盘到 page cache、CPU page cache 到用户缓冲、CPU 用户缓冲到 socket 缓冲、DMA socket 缓冲到网卡)和四次上下文切换;FileChannel.transferTo 在满足条件时走 sendfile(2),数据不进用户态。本篇说明 broker 中的对应路径。
consumer 拉取消息时,Kafka broker 在满足条件时用 FileChannel.transferTo 把日志段文件的数据从 page cache 直接交给 socket 发送缓冲区,再由 DMA 传到网卡。整条路径中 broker 进程不复制消息体内容,也不把消息读进 JVM 堆。前提是传输未启用 TLS:Java 的加密在用户态完成,启用 TLS 后这条路径会退化回用户态拷贝。
这条路径减少一次内存拷贝。broker 进程不参与消息体的数据搬运,CPU 不需要用于用户态的数据复制;数据始终在内核态,没有用户态和内核态之间的来回切换。Kafka 因此能在普通 JVM 堆配置下保持高吞吐。broker 只决定读哪段数据、发给谁,数据搬运交给内核。传统 read+write 需要 broker 堆容纳在途消息,会增加 GC 和拷贝开销。
批量与压缩:用延迟换吞吐
生产端 KafkaProducer 的两个参数控制攒批行为:batch.size 是单批的最大字节数,linger.ms 是消息进入批次后最多等多久就发。linger.ms=0 来一条发一条,延迟最低但吞吐最低;linger.ms>0 时 producer 会等批次凑大或超时再发,用少量延迟换更高的批量比例。消费端对称:fetch.min.bytes 是拉取响应的最小字节数,broker 端攒不够这个量就等,最多等 fetch.max.wait.ms 再返回。
// 最小可运行:攒批与压缩的关键配置
Properties p = new Properties();
p.put("bootstrap.servers", "broker1:9092,broker2:9092");
p.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
p.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
p.put("batch.size", 16384); // 单批最大 16 KB
p.put("linger.ms", 10); // 最多等 10 ms 攒批
p.put("compression.type", "lz4"); // 整批压缩
p.put("acks", "all"); // 可靠性归可靠性篇,本篇不展开
KafkaProducer<String, String> producer = new KafkaProducer<>(p);压缩和批量是配套的。compression.type 设为 snappy / lz4 / gzip 后,producer 对整个批次做一次压缩再发送;broker 默认保持压缩态(compression.type=producer,不在 broker 上解压再压),消费端拉到压缩批次后解压。压缩作用于整批,批越大压缩率越高。摊薄的是网络和磁盘的字节数:一条 1 KB 的消息单独发,网络和磁盘开销是固定的协议头加消息体;攒成 100 条的批次压缩后发,每条消息的均摊开销下降一个量级。
linger.ms 调大后,单条消息的端到端延迟上升,吞吐和压缩率也会上升。生产端攒批、消费端拉取和 broker 刷盘策略都需要在延迟和吞吐之间取舍。对于毫秒级延迟需求,需要评估这些机制带来的延迟;调参无法消除这种取舍。
分区:并行度单位与顺序性边界
分区(partition)是 Kafka 的基本并行单位。生产端按 key 分区保证同 key 消息进同一分区、分区内按写入顺序排列;消费端在一个 consumer group 内,一个分区同一时刻只分配给一个消费者,因此消费并发上限等于分区数。分区数决定了这条链路能横向扩展到多少消费者。
分区数增加会带来成本。每个分区在每个 broker 上是一组日志段文件,分区数 × broker 数 ≈ 文件句柄总量,分区过多会压文件描述符。元数据方面,controller 要维护每个分区的 leader、follower、ISR 状态,分区数增长会放大 controller 的协调开销和 leader 选举、rebalance 的时间。一般经验是单 broker 几千分区以内要观察 controller 的元数据操作延迟。
分区同时是顺序性的边界。分区内有序是 Kafka 给出的保证;跨分区没有全局顺序。需要按业务对象保序的场景用稳定分区键(如文件 ID)把同一对象的消息导到同一分区,这和 2016 第 6 篇的约定一致。要求全局有序意味着只能用一个分区,消费并发降为一,成本较高。提出「全局有序」需求时,应先确认业务是否确实需要全局有序。
边界清单
- 不适合随机读取与消息级查询。 日志是追加写、按 offset 顺序读,
.index是 offset 到文件位置的稀疏索引,只为消费定位服务;没有按消息内容或业务 key 随机查询的索引,要按内容检索只能按 offset 范围扫,或把数据同步到 ES、数据库再查。Kafka 没有为随机查询设计数据结构。 - 分区数后期调整代价高。 分区过少限制消费并发;分区过多压文件句柄、放大元数据与 rebalance 开销。分区数只能增加不能减少,而增加会改变 key 到分区的哈希路由,同一 key 的新消息落点变化,因此无法再依靠分区保证同一 key 的消息顺序;前期要按峰值消费并发预留。
- 不保证全局顺序。 分区内有序,跨分区无序。要求全局有序的场景消费并发降为一。
- 延迟与吞吐需要取舍。
linger.ms和fetch.min.bytes都会提高吞吐并增加等待时间,毫秒级延迟需求应评估其他系统。
可靠性的边界见Kafka(下):生产端、broker 与消费端的可靠性配置。acks、ISR、位移提交涉及另一组取舍,本篇不展开。broker 内部的副本同步、controller 选举与分区状态机细节,当时仅了解其存在;controller 的元数据操作与 ZK 依赖留到ZooKeeper/etcd:小数据强一致的协调服务再补。
参考资料
- Kafka 官方文档(设计章节、producer / consumer 配置,1.x 版本)
- Jay Kreps,《The Log: What every software engineer should know about real-time data’s unifying abstraction》(2014,LinkedIn 时期设计文章)
- 流对流的失败率、NIO 与零拷贝(2016 系列第 9 篇,IO 拷贝与 sendfile 分析,本篇在 broker 层面呼应)
- Kafka 异步化与削峰:热点上报和最终一致性(2016 系列第 6 篇,Kafka 业务接入套路,与本篇分工)
- Kafka(下):生产端、broker 与消费端的可靠性配置(可靠性配置与边界)
