主题
Kafka 的可靠性与高性能架构
本文是消息队列系统学习系列的 L2 核心篇。前置:28. 生产者与消费者原理:批次、拉取与 Rebalance。 学完可以配合面试题食用:Kafka 高性能核心设计、Kafka Exactly-Once 语义
Kafka 靠什么同时做到高吞吐又不丢数据
Kafka 的吞吐量在各大 MQ 里排第一,单机能做到百万级消息/秒。同时它又宣称"不丢消息"(配合 acks=-1 配置)。这两个目标听起来矛盾——高吞吐往往意味着牺牲可靠性,Kafka 怎么同时做到?
答案是:Kafka 在写入路径做了大量优化(顺序写、零拷贝、批量压缩),在可靠性上依赖一套副本+ISR 机制,两条线独立设计,不互相挤兑。下面逐个拆开看。
高性能三板斧
1. 顺序写 + 页缓存
Kafka 的消息写入 Broker 时,不会对磁盘做随机写,而是把所有消息追加到 Partition 对应的日志文件末尾。顺序写磁盘的速度(约 600 MB/s)比随机写(约 0.1 MB/s)快三个数量级。
但 Kafka 不只是顺序写——它写的是页缓存(Page Cache),不是直接写磁盘。操作系统管理着内存中的页缓存,Kafka 把消息塞进页缓存就返回,刷盘是 OS 异步干的。这和 Redis 的 AOF 追加写类似,但 Kafka 利用了 OS 的缓存策略,不需要自己实现一套缓存管理。
读的时候呢?消费者大概率直接从页缓存命中,不需要磁盘 I/O。这就是"读写分离"的简化版本:写更新页缓存,读也优先从页缓存取。
2. 零拷贝(sendfile)
消费者拉取消息时,传统做法需要把数据从磁盘读到内核空间,再拷贝到用户空间,再通过 Socket 发出去。这中间经过 4 次上下文切换和 4 次数据拷贝。
Kafka 用 sendfile 系统调用,让数据直接从磁盘文件(或页缓存)通过 DMA 拷贝到网卡,经过内核空间但不经过用户空间。数据拷贝次数从 4 次降到 2-3 次,上下文切换从 4 次降到 2 次。
解释一下为什么 Kafka 能用零拷贝而其他程序不一定能用:零拷贝要求数据在传输过程中不需要修改元数据。Kafka 的消息格式固定,不需要在读取时做序列化/反序列化,天然适合这个模式。
3. 批量 + 压缩
生产者不会发一条消息就做一次网络请求。客户端把消息攒到一批(batch),达到 batch.size 或 linger.ms 阈值才发出去。一批消息一次性压缩(gzip/snappy/lz4/zstd 可选),压缩比取决于消息体大小和重复度——文本日志可以压到 1/10 甚至更低。
批量发送不仅减少网络开销,还降低了 Broker 端 I/O 压力。一条 1 KB 的消息,一万条就是 10 MB,压缩后可能只有 2 MB,磁盘写和网络传输都省了。
可靠性:副本、ISR、HW/LEO
ISR 机制
Kafka 的每个 Partition 有多个副本,其中一个是 Leader,其余是 Follower。Leader 负责读写,Follower 从 Leader 拉取数据保持同步。
ISR(In-Sync Replicas)是"与 Leader 保持同步的副本集合"。什么叫"保持同步"?Follower 在 replica.lag.time.max.ms(默认 30 秒)内没有落后 Leader 太多(这里的"太多"从 Kafka 0.9 起改为时间阈值,不按消息条数算)。超时的 Follower 被踢出 ISR,等它追上再重新加入。
HW 和 LEO
- LEO(Log End Offset):副本当前最后一条消息的 offset
- HW(High Watermark):ISR 中所有副本都确认到的最小 offset。消费者只能读取 HW 之前的消息,HW 之后的消息即使 Leader 已经写入,也视为"未确认"
流程(以 acks=all 为例):
- Producer 发消息到 Leader,Leader 写入本地日志,LEO 推进
- Follower 拉取消息,写入本地,返回确认给 Leader
- Leader 发现 ISR 中所有副本的 LEO 都 >= 某条消息,就把 HW 推进到那个位置
- Producer 收到 Leader 的成功响应
- 消费者只能看到 HW 之前的消息
acks=-1(all)的落地含义
Producer 配置 acks=all 意味着 Leader 要等 ISR 中所有副本都确认写入后才返回成功。这是 Kafka 最可靠的写入模式。
但有个陷阱:如果 ISR 里只有 Leader 自己(其他副本都挂了),acks=all 退化为 acks=1。所以必须配合 min.insync.replicas 参数——比如设为 2,规定 ISR 至少要有 2 个副本才接受写入,否则报错。这个参数是可靠性与可用性的权衡点:设 2 意味着容忍 1 个副本故障,但 ISR 不够时生产者会报错。
Exactly-Once 语义
幂等生产者
设置 enable.idempotence=true 后,Kafka 给每条消息分配一个 Producer ID(PID)和 Sequence Number。Broker 根据 PID + Partition + SeqNum 去重,即使 Producer 重试也不会重复写入。
注意:这只保证单分区内、单会话内的不重。如果 Producer 挂了重启(PID 变了),之前的 SeqNum 状态失效,做不到跨会话幂等。
事务
Kafka 事务在幂等生产者基础上加了 __transaction_state 主题和 Transaction Coordinator,实现跨分区、跨会话的原子写入。
使用流程:
java
producer.initTransactions();
producer.beginTransaction();
// 发送消息
producer.send(record1);
producer.send(record2);
producer.commitTransaction(); // 或 abortTransaction()但事务的边界要清楚:它保证 Produce 端的原子性(要么全写要么全不写),不保证消费端的幂等。消费者如果在 offset 提交和业务处理之间挂了,重启后可能会重复消费,消费端仍需自己处理幂等。
动手实操
可靠生产者配置模板
java
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
// 可靠性配置
props.put("acks", "all"); // 等所有 ISR 副本确认
props.put("enable.idempotence", "true"); // 幂等生产者
props.put("min.insync.replicas", "2"); // 最少同步副本数
props.put("retries", Integer.MAX_VALUE); // 无限重试
props.put("max.in.flight.requests.per.connection", "5"); // 幂等模式下可设为 5
props.put("delivery.timeout.ms", "120000"); // 2 分钟超时注意:max.in.flight.requests.per.connection 在幂等模式下可以超过 1,因为 SeqNum 保证不会乱序重复。
吞吐对比测试
Kafka 自带了压力测试工具,测试命令:
bash
# 测试生产者吞吐(10 万条,每条 1KB,acks=1)
kafka-producer-perf-test \
--topic test-throughput \
--num-records 100000 \
--record-size 1024 \
--throughput -1 \
--producer-props acks=1 bootstrap.servers=localhost:9092
# 对比 acks=all 的差距
kafka-producer-perf-test \
--topic test-throughput \
--num-records 100000 \
--record-size 1024 \
--throughput -1 \
--producer-props acks=all bootstrap.servers=localhost:9092同一台机器上,acks=1 和 acks=all 的吞吐差距通常在 10%-30% 之间,取决于副本数和 ISR 大小。这个差距说明:Kafka 的可靠性代价并不大,因为顺序写和批量压缩是主要优化手段,等待副本确认的延迟相对较小。
分区数怎么定
分区数不是越多越好。每增加一个分区,Broker 就要多维护一套元数据、文件句柄和日志段。常见的三个约束:
- 吞吐目标:分区数 × 单个分区的读写速度 = 集群总吞吐。如果目标 10 万 TPS,单分区能跑 1 万 TPS,至少需要 10 个分区(算上副本冗余)
- 延迟:分区越多,消息在 Broker 端的分发和副本同步开销越大,延迟会略微上升
- 文件句柄:每个分区对应多个日志段文件,Broker 上维护的文件句柄数 ≈ 分区数 × 副本因子 × 2。Linux 默认文件句柄限制是 1024,需要调大
实践中,分区的上下限建议:Topic 分区数 = 集群总 Broker 数 × 2 到 4 倍,作为初始值。不宜超过 1000(单 Broker 不超过 200)。
常见误区与小结
- 误区:零拷贝 = 数据不经过磁盘。零拷贝只是在数据从磁盘到网卡的路上省去了 CPU 拷贝,磁盘 I/O 还是发生了。如果数据在页缓存里,确实不读磁盘,但那是页缓存的功劳,不是零拷贝给的。
- 误区:acks=all 保证不丢消息。严格说,它保证 ISR 中所有副本都确认了,但如果 ISR 全挂,数据可能丢失。多副本 + 定时刷盘才是完整方案。
- 误区:分区数越多吞吐越高。分区数存在拐点,超过某个阈值后,Broker 的元数据维护和文件句柄开销会反噬吞吐。
- 误区:Kafka 事务能保证消费端不重复。事务只保证 Produce 端原子性,不保证消费端幂等。
- 误区:ISR 中所有副本都同步后才响应。实际上 Follower 可以异步拉取,只要在
replica.lag.time.max.ms内跟上就行,不是真正的同步复制。
小结:Kafka 的高性能来自顺序写、页缓存、零拷贝和批量压缩,四者互不依赖,独立生效。可靠性由 ISR + HW/LEO + acks 机制保证,与高性能路径不冲突。幂等生产者和事务提供了 exactly-once 写入能力,但消费端仍需自行幂等。下一篇进入 30. RocketMQ 详解:架构、事务消息与延时消息,对比 Kafka 和 RocketMQ 的设计差异。
参考
参考:Kafka 官方文档:Exactly Once Semantics、Kafka 官方文档:Replication、《Kafka: The Definitive Guide》第 4-5 章