Skip to content

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.sizelinger.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 为例):

  1. Producer 发消息到 Leader,Leader 写入本地日志,LEO 推进
  2. Follower 拉取消息,写入本地,返回确认给 Leader
  3. Leader 发现 ISR 中所有副本的 LEO 都 >= 某条消息,就把 HW 推进到那个位置
  4. Producer 收到 Leader 的成功响应
  5. 消费者只能看到 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=1acks=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 SemanticsKafka 官方文档:Replication、《Kafka: The Definitive Guide》第 4-5 章

手撕 → 框架 → 生产化,一步步把 AI Agent 工程化搞透。
粤ICP备2026104257号-1