我最早意识到分区策略不能随便配,是在一个订单系统上。当时 Kafka 集群只开了 3 个分区,所有订单消息都按用户手机号哈希路由,看似合理,结果晚上促销一开,某几个大客户的单量直接打满一个分区,下游消费者处理不过来,消息堆积肉眼可见。从那以后,我对分区策略的理解就从“选个数”变成了“路由规则、并发模型、顺序保证三者的平衡”。
这篇文章想把 Kafka 分区策略这条线完整梳一遍:生产者怎么把消息放进不同分区、消费者怎么从分区拉数据、分区数量怎么定、自定义 Partitioner 怎么写、线上遇到的问题怎么排查。适合刚接触 Kafka 的同学建立整体认知,也适合正在调优或准备面试的开发者参考,内容全部基于我在生产环境里的实操经验,尽量说大白话。
1. 先把分区聊透:分区策略解决的核心问题
1.1 分区是并行度的基石,不是“多存几份数据”
很多刚开始学 Kafka 的人会把分区和副本搞混。分区是一份数据按路由规则拆成的多个子集,每个分区在同一个消费者组里只能被一个消费者线程消费,分区越多并行度越高。副本才是用来做数据冗余的,比如一个分区有 3 个副本,它们之间是主备关系,并不增加消费并行度。
可以用高铁站来类比:分区就像多个检票口,检票口越多,单位时间能通过的人越多;副本就像每个检票口安排了候补人员,一个倒下了另一个顶上来。分区策略则是决定你走向哪个检票口的规则,而这个规则恰恰是吞吐量和消息顺序的命门。
为什么 Kafka 不用一条大队列扛所有消息?因为单队列的瓶颈在顺序写入和单节点吞吐。Kafka 把 Topic 拆成多个分区后,这些分区可以分布在集群的不同 broker 上,写到不同磁盘,水平扩展能力一下子就出来了。Topic 的理想吞吐上限约等于“单分区吞吐 × 分区数”,所以分区数量直接影响系统的并发天花板。消费端也一样,消费者组里最多能同时跑多少个活跃消费者,完全由主题被分配到的分区数决定,消费者数量超过分区数时多出来的消费者只能空等。
1.2 一条消息从发送到落盘,到底经历了什么
Producer 发送一条消息,内部并不是直接发给 broker。完整链路是:客户端先对消息的 key 和 value 做序列化,然后交给分区器 Partitioner,由它选出一个目标分区号;消息进入 RecordAccumulator 中该分区对应的批次队列,凑满一个 batch(或者到达 linger.ms 超时)后,Sender 线程才把这一批数据发送给该分区 leader 副本所在的 broker;broker 写入日志,follower 从 leader 同步数据,完成副本复制。
这里要特别强调一个很多人忽略的点:决定消息进哪个分区的,是客户端的分区器,而不是 broker。这意味着你完全可以通过实现 Partitioner 接口,把业务路由规则写进发送端逻辑里。也就是说,分区策略的可定制性很强,不一定非得用 Kafka 默认的轮询或者哈希。
默认情况下,Kafka 客户端的处理逻辑是:如果消息带 key,使用哈希算法对 key 取哈希值,再对分区总数取模,得到分区号;如果消息不带 key,老版本采用轮询方式,新版本(2.4 之后)采用粘性分区策略。无论哪种,核心目的都是让消息尽量均匀地分布在各个分区,同时兼顾批次效率。理解了这条链路,再去聊各种分区策略,思路就顺了。
2. 常见分区策略拆解:从默认轮询到按 key 路由
2.1 无 key 场景:轮询、随机与粘性分区的取舍
如果业务不关心顺序,消息也没有明显归属维度,不带 key 是最常见的用法。此时默认分区器会尽量把消息摊到所有分区上去。
早年的默认策略是轮询(Round-Robin)或者随机(Random),每条消息依次落到不同分区。这种做法的优点是分布均匀,缺点也很明显:每条消息都可能进入一个新的 batch 窗口,批次容易被拆散,导致发送请求数变多,吞吐上不去。Kafka 2.4 起默认改为粘性分区(Sticky Partitioner),思路是:当前批次还没填满之前,连续的消息都往同一个分区写,等这个批次的 buffer 写满或者达到 linger.ms 超时,再“粘”到下一个分区继续写。
粘性分区的收益非常好理解:攒批次数变多,单次发送的数据包变大,请求次数明显减少,broker 端的压力也跟着变小。有测试数据显示,在同样的 qps 下,粘性分区比轮询能减少 30% 以上的请求量,这对高吞吐场景是很可观的优化。所以如果你没有特殊需求,无 key 场景建议直接交给默认策略,不需要自己再造轮子。只是在脑中对这个特性有数,以后看监控发现短时间内消息集中落在某一个分区,别慌,这大概率是粘性分区的正常表现。
2.2 有 key 场景:哈希路由与顺序性收益
按 key 哈希路由是 Kafka 里最常用的分区策略,也是保证消息顺序的基本手段。思路很简单:相同 key 的消息会经过相同的哈希计算,最终一定进入同一个分区。同一个分区内部的消息是有序的,消费者按顺序拉取,就能实现“key 级别的顺序保证”。
拿电商订单场景举例。订单创建、订单支付、订单退款这些消息都携带同一个订单号作为 key,它们被路由到同一个分区,消费者端按照顺序处理,就能保证“先创建、再支付、后退款”的业务逻辑不乱套。如果不做这个设计,两条关联消息落到不同分区,被不同消费者线程并行处理,顺序就完全不可控了。
这里要补一个细节:Kafka 默认的哈希不是 Java 的 hashCode(),而是 murmur2 哈希。原因是 Java hashCode 的分布质量和不同 JVM 版本之间的稳定性都有隐患,murmur2 的散列性更好,也更稳定。以前见过有人自己实现 key.hashCode() % numPartitions,一旦换 JDK 版本或者字符串哈希算法调整,分区路由就乱掉,这是典型的想当然坑。
2.3 分区数变化带来的顺序性破坏
按 key 哈希有一个很隐蔽的风险点:分区数量不是一成不变的。Kafka 允许增加分区数,但明确禁止减少分区数。一旦分区数从 N 变成 M,同一 key 的模值可能随之变化,新消息会进入新的分区,而历史消息还在老分区里躺着,双方的处理进度不一致,顺序自然就保不住了。
举个例子:某个订单在 10 个分区时被路由到分区 0,消息延迟了 20 分钟才投递成功,期间运维把分区扩容到 20 个,后续关联消息被路由到分区 11。消费端如果并行消费这 20 个分区,两条关联消息就可能在处理时“错位”。这种问题非常难排查,因为日志里每条消息本身都是对的,只是业务侧先看到了后发生的消息。
所以在生产环境里,Topic 分区数一定要在业务上线前定好,并留足余量。不能想着“先建 3 个分区跑起来,后面不够再加”,一旦业务开始对顺序有要求,后期加分区就是一次线上事故。如果确实需要扩容,要先评估下游消费逻辑是否依赖全局或 key 级顺序,必要时采用双写迁移或者短期停写方案。
2.4 自定义分区器:把业务路由规则写进 Partition 环节
自定义分区器的使用场景,除了应对热点 key 之外,还包括多租户隔离、数据本地性、按时间分片、权重分配等。比如多租户场景下,你希望 A 租户的消息稳定落在分区 0~5,B 租户落在分区 6~11,方便按租户做流量隔离或配额管理。这属于典型的自定义分区诉求。
写自定义分区器有几个设计红线需要记住:第一,partition() 方法在生产发送链路上会被高频调用,内部绝对不能出现 IO 操作、远程调用、数据库查询这类耗时逻辑,否则发送性能直接崩塌。第二,key 为 null 时必须有兜底处理,不要抛 NullPointerException。第三,如果返回的分区号超出真实分区数,客户端会报“Invalid partition”错误,所以最好在内部做一次取模约束。第四,要考虑未来分区数变化时的兼容性,别把分区数硬编码在业务逻辑里。
至于实现细节和完整代码,我会在第 4 章专门演示,那里会给出一个可以直接搬到项目里的 Partitioner 写法。
3. 生产端只是半程:消费端如何分配分区
3.1 消费者组再平衡与分区分配策略
生产者把消息写进了分区,消费端还需要决定“哪个消费者来消费哪个分区”。这个决策过程由消费者组和 GroupCoordinator 协作完成。消费者组里的成员会定期发送心跳,一旦有消费者加入、退出、订阅 Topic 变化,或者分区数变化,就会触发再平衡(Rebalance),把分区重新分配给组内成员。
再平衡期间,整个消费者组会短暂停止消费,所以 Rebalance 的频率和耗时直接影响消费稳定性。Kafka 提供了多种分区分配策略,默认的在较新版本里是 CooperativeStickyAssignor(KIP-429 之后,大概 3.1 版本开始成为默认),之前的老版本默认是 RangeAssignor。
不同策略的区别在于分配的动作:
- RangeAssignor:按 Topic 逐个分配,每个 Topic 的分区先除以消费者数,余数分给靠前的消费者。问题是有可能不均匀,尤其当多个 Topic 的分区数、消费者数不成比例时。
- RoundRobinAssignor:把订阅的所有分区放在一起轮询,要求组内所有消费者订阅的 Topic 列表一致,否则分配会失衡。
- StickyAssignor / CooperativeStickyAssignor:尽量保持已有分配不变,一次 Rebalance 只移动必要分区的归属,减少分区在消费者之间的“搬家”开销。
如果你还在用老版本客户端,或者遇到了频繁 Rebalance 导致消费停滞的问题,建议检查一下分配策略配置,配合 session.timeout.ms 和 heartbeat.interval.ms 把稳定性调起来。
3.2 分区数、消费者数与吞吐量的三角关系
分区数是 Kafka 吞吐能力的关键变量,但它并不是越大越好。分区数和消费者数的关系跟“并发度上限”强绑定:同一个消费者组内,一个分区同时只能被一个消费者线程消费,所以最大有效消费者数等于被分配到的分区总数。消费者数量超过分区数,多出来的消费者就是空转,纯属浪费资源;分区数远大于消费者数,单个消费者要处理多个分区,吞吐瓶颈在单个消费者的处理能力上。
我在实际项目里常用的分区数估算方法很简单:先估计业务峰值时的单分区处理能力,然后按目标吞吐反推。比如单分区每秒能处理 2000 条消息,业务峰值是 40000 条每秒,那至少需要 20 个分区。再留 30%~50% 的缓冲,最终定在 28~30 个分区。另外还有一类经验值是“分区数不超过 broker 数的 10 倍”,比如 3 台 broker 的集群,Topic 分区数控制在 30 以内比较稳妥,否则文件句柄和 Rebalance 的代价都会变大。
这里顺带提一句流计算场景。Flink 消费 Kafka 时,KafkaSource 的并行度一旦超过分区数,多余的任务就会闲置;如果并行度小于分区数,又有多个分区被同一个任务消费的问题。所以通常会把 Flink 并行度设置为分区数的整数倍,这样既能打满分区,又方便下游算子做并发汇聚。
3.3 顺序消费的完整链路设计
很多面试题会问“Kafka 怎么保证消息有序”,但真实答案是:Kafka 默认不保证全局有序,只保证分区内有序。想要实现业务上的顺序,需要全链路配合。
第一步,生产端按 key 哈希,保证同一业务的关联消息进同一分区;第二步,消费端该分区的消息由一个消费者线程顺序拉取、顺序处理;第三步,如果消费者内部做了异步并发处理,要自己解决乱序问题,比如按业务 ID 做分桶并发、处理完按序号回写。这三步缺一环,顺序就保不住。
还有一类比较常见的需求是“延迟处理”。热搜词里那个“kafka 如何延迟30分钟消费”,在电商里就是下单 30 分钟未支付自动关闭。Kafka 原生没有延迟队列能力,常见的做法是消费者收到消息后先写入 Redis ZSet,score 设置为期望执行的时间戳,另起一个定时任务每分钟扫一次,到期再把消息重新投递到业务 Topic;或者直接用 RabbitMQ 延迟插件、RocketMQ 定时消息,比硬刚 Kafka 简单得多。
有些同学想通过 pause()/resume() 让 Kafka Consumer 暂停多少分钟再继续,这个做法会让整个消费流程卡住,对同 Topic 的其他消息也有影响,不推荐在生产环境用。
4. 代码实操与生产环境排障实录
4.1 自定义分区器的完整代码与配置参数
写一个按租户 ID 分区的 Partitioner。假设 key 格式是 {tenantId}:{bizId},tenantId 为 3 位字符串,比如 “001:ORDER123456”。业务要求同一租户的消息尽量落在一起,并且给大租户预留更多的分区。代码如下:
import org.apache.kafka.clients.producer.Partitioner; import org.apache.kafka.common.Cluster; import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.record.InvalidRecordException; import java.util.List; import java.util.Map; public class TenantPartitioner implements Partitioner { // 大租户,可以按需调整 private static final String LARGE_TENANT = "007"; @Override public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) { List<PartitionInfo> partitions = cluster.partitionsForTopic(topic); int numPartitions = partitions.size(); if (key == null) { // key 为 null 时兜底,走随机或轮询,这里简单取当前时间戳 int random = (int) (System.currentTimeMillis() & 0x7fffffff) % numPartitions; return random; } String keyStr = key.toString(); if (!keyStr.contains(":")) { throw new InvalidRecordException("key format error: " + keyStr); } String tenantId = keyStr.split(":")[0]; // 大租户独占后半部分分区 int half = Math.max(1, numPartitions / 2); if (LARGE_TENANT.equals(tenantId)) { return half + (tenantId.hashCode() & 0x7fffffff) % (numPartitions - half); } // 普通租户集中在前面分区 return (tenantId.hashCode() & 0x7fffffff) % half; } @Override public void close() { } @Override public void configure(Map<String, ?> configs) { } }这个例子最核心的地方在于,它演示了“按业务规则切分分区段”。大租户只占总租户数的一小部分,但流量可能占 80%,单独划走一半分区,避免普通租户的消息被大租户的洪峰冲垮。普通租户哈希到前半段,大租户哈希到后半段,两者互不干扰。
配置方式有两种。Spring Boot 的 application.yml 里这样写:
spring: kafka: producer: properties: partitioner.class: com.example.demo.TenantPartitioner原生 Kafka API 这样写:
props.put(ProducerConfig.PARTITIONER_CLASS_CONFIG, TenantPartitioner.class.getName());配置完成后,可以先用控制台生产者发几条不同前缀 key 的消息,再通过可视化工具看分区分布,确认路由是否符合预期。这里提醒一句:不要在生产环境频繁切换分区器实现,路由规则一旦变了,存量消息和增量消息会分到不同分区,顺序性和数据分布都会受影响。
4.2 可视化工具排查分区分布:Offset Explorer 连接与常见报错
排查分区分布和消费 lag,我一般用 Offset Explorer(原名 Kafka Tool),免费版就够用。它能看到 Topic 的每个分区、leader 副本、ISR 列表、消息条数以及消费者组的消费进度,信息很直观。
连接本地单机 Kafka 的配置很简单:在 Offset Explorer 里新增集群,Properties 面板选择 Bootstrap servers,填 localhost:9092,其他保持默认,即可读取到集群信息。如果是 Docker 容器里起的 Kafka,记得不要用容器内部主机名,宿主机访问地址要配成 localhost:9092 或者宿主机的局域网 IP。ADVERTISED_LISTENERS 配置不对,外部工具连不上是最常见的问题。
再来说几个高频报错。Error while fetching metadata with correlation id 这个提示,本质是客户端连不上 bootstrap.servers 或者连接后拿不到元数据,原因可能是 broker 地址不对、ACL 没授权、集群没起来。Cluster authorization failed 则基本是 ACL 权限问题,往消费者组前缀和 Topic 前缀对应的 producer/consumer 角色上加权限即可。Timed out 则要重点检查地址和网络,比如容器内用 localhost 访问不了宿主机上的 broker,要改成宿主机 IP。
除了 Offset Explorer,开源的 Kafka UI、kafdrop 也可以看分区和消费组,各有侧重。我个人的习惯是:本地调试用 Offset Explorer,集群上想快速看消息内容用 Kafka UI,功能比较全。
4.3 Docker 环境跑通 Kafka 集群的注意事项
新版 Kafka 从 3.x 开始进入了 KRaft 模式,可以完全不依赖 Zookeeper。很多人一开始不熟悉这个模式,README 上也找不到 zookeeper 配置,容易卡住。下面给一个最简示例,用 apache/kafka 官方镜像直接启动:
docker run -d --name kafka \ -p 9092:9092 \ -e KAFKA_PROCESS_ROLES=broker,controller \ -e KAFKA_NODE_ID=1 \ -e KAFKA_CONTROLLER_QUORUM_VOTERS=1@localhost:9093 \ -e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093 \ -e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092 \ -e KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=PLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT \ -e KAFKA_CONTROLLER_LISTENER_NAMES=CONTROLLER \ apache/kafka:3.7.0这个命令里 KAFKA_PROCESS_ROLES 声明当前节点同时担任 broker 和 controller,单节点不需要再单独起一个 controller 容器。KAFKA_CONTROLLER_QUORUM_VOTERS 配置的是 controller 节点列表,格式是 nodeId@host:port。KAFKA_ADVERTISED_LISTENERS 是客户端真正用来连接 broker 的地址,这也是最容易踩坑的地方,默认如果写成 kafka:9092,宿主机上的客户端必然连不上。
多节点集群部署时,每个 broker 的 KAFKA_NODE_ID 要不同,CONTROLLER_QUORUM_VOTERS 要把所有 controller 节点写全,并且 controller 端口要暴露出来。另外我不会建议在 Docker 里跑生产级集群,容器重启后 IP 变化会导致元数据里的 broker 地址失效,本地学习或者测试用是可以的,生产还是老老实实用裸机或 K8s 有状态服务。
5. 高频面试题与避坑速查表
5.1 Kafka 分区相关问题速查表
无论是准备面试还是自查,下面这些点都是绕不开的,我做成了速查表方便对照:
| 问题 | 核心答案要点 |
|---|---|
| Kafka 为什么吞吐高 | 分区并行 + 顺序写 + 页缓存 + 零拷贝 |
| 分区数是否越多越好 | 不是,受文件句柄、Rebalance、controller 压力限制 |
| 分区数能减少吗 | 不能,Kafka 只支持增加分区 |
| 怎么保证消息有序 | 生产者按 key 哈希保证同 key 进同分区,消费者单线程处理,不异步并发 |
| 分区分配策略有哪些 | Range、RoundRobin、Sticky、CooperativeSticky |
| 消费者数超过分区数会怎样 | 多余消费者空转,不提升吞吐 |
| 自定义分区器要注意什么 | partition() 内避免 IO、key 为 null 要兜底、分区号不能越界 |
| Topic 分区数和 broker 数关系 | 经验建议分区总数不超过 broker 数的 10 倍左右(看机器规格) |
5.2 我在生产环境踩过最狠的几个坑
第一个坑,分区数量定得过于随意。早期一个日志类 Topic 直接建了 60 个分区,集群只有 3 台 broker,单 Topic 的文件句柄和副本同步压力都很大,后来高峰时段的 Rebalance 时间从几秒飙到几十秒,整个消费链路都在等分区重分配。从那以后我养成了估算的习惯,同时也要结合下游消费能力来定,分区不是为 Kafka 面子定的,是为消费并发定的。
第二个坑,自定义分区器里加了远程 Redis 查询。当时我想按会员等级分流,直接在 partition() 里实时查 Redis 判断大客户,结果生产一上流量,producer 吞吐直接掉了一个数量级。后来改成客户端启动时加载等级映射表,定时刷新,分区器里只做纯内存计算,问题才解决。记住我前面那条红线:partition() 是热路径,绝对不能碰远程 IO。
第三个坑,消费者组频繁 Rebalance。线上消费者日志一直在刷 rebalance,消费进度走走停停。排查后发现是 session.timeout.ms 设置太小,消费者偶尔 GC 停顿超过阈值就被踢出组,触发新一轮分配。调大 session.timeout、加大 heartbeat.interval 的宽容度、顺便把消费者端的 max.poll.interval.ms 调大之后,Rebalance 频率才算压下去。如果你也遇到类似问题,建议先看日志里是哪种 Rebalance 原因,再对症下药。
第四个坑,使用 Spring Boot 多 Kafka 地址消费时,配置串了。Spring Boot 对多个 KafkaTemplate 和消费者工厂的管理容易混淆,下游用了错误的 bootstrap.servers,报错信息却指向“cluster authorization failed”,排查半天才发现是连到了测试环境集群。多环境多集群场景下,建议每个集群单独命名一个消费者工厂,bootstrap.servers 单独从配置中心注入,别混用一套默认配置。
5.3 分区相关问题的日常监控手段
分区策略上线的效果不能全靠感觉,监控指标要跟上。最核心的几个指标:每个分区消息堆积量(Lag)、消费者组 Rebalance 次数和时间、Topic 分区消息不均度、Producer 的请求速率和 batch 大小。
Lag 可以用 Kafka 自带命令 kafka-consumer-groups.sh 定期拉取,也可以在监控面板上集成。如果发现某个分区 Lag 长期高于其他分区,大概率是分区器路由不均匀或者消费者处理某类消息更慢,需要进一步看 key 的分布。Rebalance 次数要纳入告警,每分钟超过 1 次就该拉响警报。Producer 端请求速率如果明显下降,看看 batch 是否被频繁切换,可能需要调大 linger.ms 或 batch.size。
我个人比较推荐用 Prometheus + Grafana 那一套,Kafka 官方有 exporter,加上 consumer lag 的 exporter,基本能覆盖分区和消费链路的主要指标。把这些指标盯住了,很多分区策略的问题能在影响业务之前就暴露出来。
我个人的体会是,分区策略从来不是一锤子买卖,它是“路由规则—分区数—消费模型”三者咬合的系统设计。动手前先把业务流量模型摸清楚,再决定怎么分区;上线后要持续盯分区堆积、消费者 Lag 和 Rebalance 频率;改分区数前先确认下游没有依赖旧分区的顺序假设。最后分享一个小技巧:如果某个 key 的流量异常大,可以在 key 后面拼一个随机后缀再路由,同时把原始 key 放在消息体里,消费端再做聚合,这样能快速削掉热点分区的压力,也是应对倾斜最简单有效的一招。