跳至主要內容

十万个why:为什么 Kafka Partition 越多写入性能越差 ?

程序员小富大约 9 分钟

大家好,我是小富~

最近在做系统吞吐量压测,团队里不知道哪个爹为提升 Kafka 的写入并发,直接在测试环境把某个 Topic 的 Partition数量从 8 个改成了 1000 个。

结果压测一跑,写入吞吐量不仅没升,反而像拉稀一样直线下滑,Broker 节点的 CPU 飙高,磁盘 I/O 更是直接拉满,甚至还报了 Too many open files 的错。

除了无语没啥好说的,一点常识没有,盲目迷信分区越多,吞吐量越高,八股文不都写了 Kafka Partition 越多,写入性能反而越差吗?

Producer 内存模型被分区数量击穿

要理解这个问题,我们得先搞清楚 Kafka Producer 写入数据的底层内存模型。

调用 producer.send() 发送消息,消息并不会直接走网络发给 Broker,而是会先进入一个叫 RecordAccumulator 的内存缓冲区。

这个缓冲区的核心数据结构是一个 ConcurrentMap,Key 是 TopicPartition,Value 是 Deque<ProducerBatch>。也就是说,每一个分区都独占一条发送队列

我们翻了一下 Kafka 源码里的 RecordAccumulator

// org.apache.kafka.clients.producer.internals.RecordAccumulator
public class RecordAccumulator {
    // 关键:每个 TopicPartition 独占一个 Deque
    private final ConcurrentMap<TopicPartition, Deque<ProducerBatch>> batches;
    // 内存池,所有分区共用这一个池子
    private final BufferPool free;
    // batch.size:每个 ProducerBatch 的目标大小,默认 16KB
    private final int batchSize;
}

每一个分区被首次写入时,RecordAccumulator 就会从共享内存池 BufferPool(由 buffer.memory 控制,默认 32MB)中申请一块 batch.size(默认 16KB)大小的 ByteBuffer 给它。

假设你只开了一个 Producer 实例,要往一个有 1000 个分区的 Topic 写数据。当这 1000 个分区全部活跃时,Producer 至少需要同时持有:1000 * 16KB = 16MB 的 ByteBuffer。

16MB 看似不大,但别忘了两个关键细节:

第一,BufferPool 的总容量默认只有 32MB。16MB 的 Batch 缓存一出去,留给数据填充的空间只剩下一半。在消息量密集时,很容易出现某些分区的 Batch 还没攒够就被新的 Batch 挤出去,导致大量半满批次被发送,网络请求数暴增,每个请求只携带极少的数据量。

第二,更关键的是 BufferPool 内部的锁竞争。看源码就知道,BufferPool.allocate() 方法内部用的是 ReentrantLock + Condition。1000 个分区意味着高频率的 allocate/deallocate 操作,大量线程在锁上排队等待,Producer 的发送线程频繁阻塞。

// org.apache.kafka.clients.producer.internals.BufferPool
public ByteBuffer allocate(int size, long maxTimeToBlockMs) throws InterruptedException {
    this.lock.lock();  // 所有分区的内存申请都要抢这一把锁!
    try {
        // 如果内存不够,就在这里阻塞等待
        if (this.nonPooledAvailableMemory + freeListSize >= size) {
            // 分配内存...
        } else {
            // 阻塞等待,直到有其他 Batch 释放内存
            this.waiters.addLast(moreMemory);
            // ...
        }
    } finally {
        this.lock.unlock();
    }
}

buffer.memory 耗尽,send() 方法会阻塞最长 max.block.ms(默认 60 秒)毫秒,之后直接抛出 TimeoutException。生产环境中,这意味着你的业务线程被卡住整整一分钟。

这就像你在家里摆了 1000 个垃圾桶,每个垃圾桶一装满就得倒,结果你家里大部分空间都被空的垃圾桶占满了,你连路都没法走。

顺序写磁盘退化成随机写

Kafka 之所以能做到单机十几万的 QPS 写入,有个主要原因是操作系统 PageCache + 顺序写

所谓的顺序写,就是每一个批次的数据在文件中连续写入,没有跳转。打个比方,就像写日记一样,一直在最后一页写,写完一页翻到下一页,几乎不用移动磁头,所以磁盘 I/O 的性能非常非常高,非常接近内存的速度。

每一个 Partition 在 Broker 的磁盘上都对应一个独立的物理文件夹。这个文件夹里至少包含:

  • .log 物理数据文件(Segment)

  • .index 偏移量索引文件

  • .timeindex 时间戳索引文件

如果在一个 Topic 里建了 1000 个分区,在物理磁盘上就会产生 1000 个文件夹,涉及至少 3000 个物理文件。Broker 密集地往这 1000 个分区刷盘,从操作系统内核的角度看,它需要在 3000 多个文件的 PageCache 之间来回切换写入位置。

在 HDD 机械硬盘上,这意味着磁头要在 1000 个不同的磁道之间来回寻址,顺序 I/O 彻底退化成了随机 I/O,吞吐量可以从 100MB/s 暴跌到个位数。

即使用的是 SSD,虽然没有磁头寻址的问题,但也有个PageCache 争用问题,Linux 内核的 PageCache 是全局共享的 LRU 缓存。3000 个文件同时竞争 PageCache 空间,会导致热数据频繁被冷数据挤出。

单台 Broker 上的分区数超过一定阈值后,即使在 SSD 上,端到端延迟 p99 也会出现 5~10 倍的劣化。

Too many open files

Linux 系统中每打开一个文件,内核就要分配一个文件描述符,Kafka 的每一个 Partition 在运行时至少需要常驻打开很多文件句柄。

普通的 Linux 服务器默认的 ulimit -n 一般是 1024。即使运维手动调高到了 65536,在分区数膨胀的场景下依然可能不够。一旦文件描述符耗尽,Broker 的表现不是优雅降级,而是灾难性的:

  • 新的 Producer/Consumer 连接无法建立

  • 新的 Segment 文件无法创建,导致写入直接报错

  • 日志文件无法写入,排障连日志都没有

java.io.IOException: Too many open files
    at sun.nio.ch.FileDispatcherImpl.open0(Native Method)
    at kafka.log.LogSegment.<init>(LogSegment.scala:xxx)

这不是慢慢变慢的问题,是直接宕机。而且一台 Broker 宕机后触发的分区迁移,会进一步加重其他 Broker 的负载,形成级联雪崩。

元数据问题

分区多了还有个问题就是 Kafka 集群的元数据太多,也会影响性能。不管用的是传统的 ZooKeeper 模式还是新的 KRaft 模式,集群中都有一个 Controller 角色负责:

  • 维护所有分区的 Leader/ISR 状态

  • 处理 Broker 上下线时的分区 Leader 选举

  • 将元数据变更广播给所有 Broker

元数据全量同步的开销

每当有分区状态变化(比如 ISR 列表收缩),Controller 需要向所有相关 Broker 发送 LeaderAndIsr 请求和 UpdateMetadata 请求。

在 ZooKeeper 模式下,每一次分区状态变更还会触发 ZK 节点的 Watch 回调。当分区数达到十万量级,一次普通的 Broker 上线/下线,Controller 需要处理的 ZK 事件可能多达数万个

Kafka 源码中 KafkaController 是单线程事件模型,内部只有一个 ControllerEventThread,所有事件排队串行处理:

// kafka.controller.KafkaController
class ControllerEventThread extends ShutdownableThread {
  override def doWork(): Unit = {
    val dequeued = queue.take()  // 从队列里逐个取事件
    dequeued.process(processor) // 串行处理!
  }
}

分区越多,这个单线程要处理的事件越密集,处理延迟越高。严重的会出现 Controller 的事件队列积压数万个未处理事件,整个集群的元数据更新拖慢了整体系统的性能。

故障恢复的选举

这是生产环境最怕的场景。一台 Broker 意外宕机,Controller 需要为该 Broker 上的每一个 Leader 分区重新选举新的 Leader。

要知道在选举完成之前,这些分区是完全不可读不可写的!!!如果你的 Topic 就只分布在两三台 Broker 上,一台挂掉意味着三分之一的分区同时不可用好几分钟,这在金融、交易等场景下是不可接受的 P0 事故。

而且选举本身还会产生大量的 LeaderAndIsrUpdateMetadata RPC 请求风暴,进一步打满网络带宽,导致存活 Broker 的正常读写也被拖慢。

Rebalance 问题

分区多了不仅影响写入,对消费端的影响也大,消费组在以下场景会触发 Rebalance 分区重新分配:

  • 消费者加入或离开消费组

  • Topic 的分区数发生变化

  • Consumer 在 max.poll.interval.ms 内没有调用 poll()

Rebalance 的耗时与分区数量正相关,原因在于:

  1. 分区分配算法的计算量:默认的 RangeAssignorCooperativeStickyAssignor 需要遍历所有分区进行分配计算。分区从 10 个变成 10000 个,分配耗时可能从毫秒级变成秒级。

  2. Stop-The-World 效应:Rebalance 期间所有消费者都必须停止消费,等待分配完成。分区越多,这个停顿窗口越长。

  3. Rebalance 的触发频率也会上升:分区多意味着单个消费者分到的分区更多,处理压力更大,更容易超过 max.poll.interval.ms 的限制,从而被踢出消费组,再次触发 Rebalance,形成恶性循环。

多少个 Partition 合理?

别去拍脑袋瞎猜,教大家一个 Kafka 官方推荐的计算公式:

分区数 = max(T/Pt, T/Ct) 其中:

  • T = 目标吞吐量
  • Pt = 单分区 Producer 吞吐量
  • Ct = 单分区 Consumer 吞吐量

测量单分区吞吐量

用 Kafka 自带的 kafka-producer-perf-test.sh 工具,在你的真实机器上压测单个分区的写入极限。通常情况下,单分区顺序写入可以达到 10MB/s ~ 50MB/s(SSD 会更高)。

# Kafka 自带的性能测试工具,测量单分区吞吐
bin/kafka-producer-perf-test.sh \
  --topic perf-test \
  --num-records 1000000 \
  --record-size 1024 \
  --throughput -1 \
  --producer-props bootstrap.servers=localhost:9092

确定目标吞吐量

假设你的业务高峰期需要 200MB/s 的写入吞吐量,单分区实测写入为 40MB/s,单分区消费为 30MB/s。

计算 max(200/40, 200/30) = max(5, 7) = 7 个分区。考虑到未来流量增长和突发峰值的冗余,乘以 1.5 ~ 2 倍,也就是 10 到 14 个分区就绰绰有余了。

绝大多数普通业务场景,单个 Topic 的分区数保持在 6 到 12 个 之间就完全能满足高并发需求,根本不需要动辄上百上千。

Kafka 官方的建议:

单台 Broker 上的总分区数(所有 Topic 的 Leader + Follower 副本之和)尽量控制在 2000 ~ 4000 个以内。

ZooKeeper 模式下,整个集群的分区总数不要超过 20 万,KRaft 模式可以适当放宽。如果用的是 HDD,这个阈值要打对折。

超过这些阈值,哪怕你用的是 NVMe SSD,Kafka 的端到端延迟和集群稳定性也会开始明显劣化。

说在最后

在 Kafka 的使用,设计 Topic 宁可先少给点分区,等不够用了再通过命令行动态增加,也绝对不要图省事一步到位建几百个分区。

因为 Kafka 支持动态增加分区,但不支持动态减少分区!一旦分区建多了,除了重新建 Topic 迁移数据,没有任何后悔药可吃。

上次编辑于: