拓冰建站拓冰建站
首页 / 资讯中心 / 正文

Kafka分区策略全解析:从路由规则到生产实践

我最早意识到分区策略不能随便配是在一个订单系统上。当时 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 副本所在的 brokerbroker 写入日志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~5B 租户落在分区 6~11方便按租户做流量隔离或配额管理。这属于典型的自定义分区诉求。写自定义分区器有几个设计红线需要记住第一partition() 方法在生产发送链路上会被高频调用内部绝对不能出现 IO 操作、远程调用、数据库查询这类耗时逻辑否则发送性能直接崩塌。第二key 为 null 时必须有兜底处理不要抛 NullPointerException。第三如果返回的分区号超出真实分区数客户端会报“Invalid partition”错误所以最好在内部做一次取模约束。第四要考虑未来分区数变化时的兼容性别把分区数硬编码在业务逻辑里。至于实现细节和完整代码我会在第 4 章专门演示那里会给出一个可以直接搬到项目里的 Partitioner 写法。3. 生产端只是半程消费端如何分配分区3.1 消费者组再平衡与分区分配策略生产者把消息写进了分区消费端还需要决定“哪个消费者来消费哪个分区”。这个决策过程由消费者组和 GroupCoordinator 协作完成。消费者组里的成员会定期发送心跳一旦有消费者加入、退出、订阅 Topic 变化或者分区数变化就会触发再平衡Rebalance把分区重新分配给组内成员。再平衡期间整个消费者组会短暂停止消费所以 Rebalance 的频率和耗时直接影响消费稳定性。Kafka 提供了多种分区分配策略默认的在较新版本里是 CooperativeStickyAssignorKIP-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 ZSetscore 设置为期望执行的时间戳另起一个定时任务每分钟扫一次到期再把消息重新投递到业务 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) { ListPartitionInfo 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(MapString, ? 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_ROLESbroker,controller \ -e KAFKA_NODE_ID1 \ -e KAFKA_CONTROLLER_QUORUM_VOTERS1localhost:9093 \ -e KAFKA_LISTENERSPLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093 \ -e KAFKA_ADVERTISED_LISTENERSPLAINTEXT://localhost:9092 \ -e KAFKA_LISTENER_SECURITY_PROTOCOL_MAPPLAINTEXT:PLAINTEXT,CONTROLLER:PLAINTEXT \ -e KAFKA_CONTROLLER_LISTENER_NAMESCONTROLLER \ apache/kafka:3.7.0这个命令里 KAFKA_PROCESS_ROLES 声明当前节点同时担任 broker 和 controller单节点不需要再单独起一个 controller 容器。KAFKA_CONTROLLER_QUORUM_VOTERS 配置的是 controller 节点列表格式是 nodeIdhost: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 放在消息体里消费端再做聚合这样能快速削掉热点分区的压力也是应对倾斜最简单有效的一招。
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门