Kafka 消息始终发送到分区 0 的原因与解决方案

云浩吖_7732

云浩吖_7732

2026-09-29

125人浏览

原创

Kafka 消息始终发送到分区 0 的原因与解决方案

kafka 消息即使设置了不同 key,仍全部路由至 partition 0,根本原因是 kafka 3.3+ 默认启用了 kip-794 引入的“严格均匀粘性分区器(strictly uniform sticky partitioner)”,它会将同一批次(batch)内的所有消息强制分配到同一分区,而非按 key 哈希独立计算。

kafka 消息即使设置了不同 key,仍全部路由至 partition 0,根本原因是 kafka 3.3+ 默认启用了 kip-794 引入的“严格均匀粘性分区器(strictly uniform sticky partitioner)”,它会将同一批次(batch)内的所有消息强制分配到同一分区,而非按 key 哈希独立计算。

在您的代码中,虽然为 pushDataRequestChannel 和 processDataRequestChannel 分别配置了不同的静态 key(如 "group_id" 和 "partition_1_key"),并期望它们分别哈希到 partition 0 和 partition 1(因 2186850892 % 2 == 0,1550936367 % 2 == 1),但实际行为受 Kafka 客户端默认分区器策略支配——并非 key 决定分区,而是批次粘性优先。

自 Kafka 3.3 起(对应 kafka-clients >= 3.3.0),DefaultPartitioner 已升级为 UniformStickyPartitioner(KIP-794),其核心逻辑是:

  • 若消息无 key(key == null),则使用粘性分区(sticky partition):为每个 topic 维护一个“当前活跃分区”,同 batch 内所有无 key 消息均发往该分区;
  • 若消息有 key,则仍按 Murmur2 哈希 + 取模计算分区(即 hash(key) % numPartitions) —— 但关键限制在于:当多条带 key 的消息被快速连续发送、且未触发立即发送(即未填满 batch.size 或未超时 linger.ms)时,它们可能被攒批(batched)进同一个 ProducerBatch;而该 batch 一旦选定首个消息的分区(基于其 key),后续同 batch 内所有消息(无论 key 是否不同)都将强制路由至该分区,以提升压缩效率和吞吐。

这正是您观察到“所有消息都进 partition 0”的原因:

  • 您的两个 MessageHandler 实例共用同一个 KafkaTemplate(即共享底层 KafkaProducer);
  • 在循环中快速发送(无显式延时),导致 pushDataRequestChannel 和 processDataRequestChannel 发出的消息被合并进同一 ProducerBatch;
  • 第一条消息(例如 i=0,走 pushDataRequestChannel,key="group_id" → hash%2=0)决定了整个 batch 的目标分区为 0;
  • 后续消息(i=1,3,5… 使用 "partition_1_key")虽 key 不同,但仍被“粘”在 partition 0。

✅ 解决方案如下:

1. 显式禁用粘性分区(推荐用于 key 驱动场景)
在 application.yml 或 KafkaTemplate 配置中设置:

spring:
  kafka:
    producer:
      properties:
        partitioner.class: org.apache.kafka.clients.producer.internals.DefaultPartitioner
        # 注意:Kafka 3.3+ 中 DefaultPartitioner 即 UniformStickyPartitioner,
        # 但可通过以下参数关闭粘性行为
        partitioner.ignore.keys: false  # 确保 key 生效(默认 true 表示忽略 key 用粘性)

更可靠的方式是降级为经典分区器(适用于 Kafka ≥ 3.3):

@Bean
public KafkaTemplate<string string> kafkaTemplate(ProducerFactory<string string> factory) {
    KafkaTemplate<string string> template = new KafkaTemplate(factory);
    // 强制使用旧版分区逻辑(非粘性、纯 key 哈希)
    template.setProducerListener(new LoggingProducerListener());
    return template;
}</string></string></string>

并在 ProducerFactory 中注入自定义 DefaultPartitioner(需 Kafka

@Bean
public ProducerFactory<string string> producerFactory() {
    Map<string object> props = new HashMap();
    props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
    props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
    props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
    // 关键:禁用粘性,确保 key 哈希生效
    props.put(ProducerConfig.PARTITIONER_CLASS_CONFIG, "org.apache.kafka.clients.producer.internals.DefaultPartitioner");
    props.put("partitioner.ignore.keys", "false"); // 必须设为 false!
    return new DefaultKafkaProducerFactory(props);
}</string></string>

2. 强制刷新批次(调试用,不推荐生产)
在每次 send() 后调用 flush(),避免攒批:

kafkaTemplate.send(topic, key, value).get(); // 同步发送确保落盘
kafkaTemplate.flush(); // 强制清空当前 batch

⚠️ 注意事项:

  • flush() 会显著降低吞吐,仅用于验证逻辑;
  • 确保 key 字符串编码一致(如 UTF-8),避免哈希值偏差;
  • 使用 kafka-topics.sh --describe 验证 topic 分区数确为 2;
  • 可通过日志开启 org.apache.kafka.clients.producer.internals DEBUG 级别,观察分区选择过程。

总结:Kafka 的“一致性哈希路由”前提,是消息能被独立评估分区——而粘性分区器通过批次优化牺牲了该确定性。理解 KIP-794 的设计权衡,并根据业务需求(key 敏感型 or 吞吐优先型)合理配置 partitioner.ignore.keys 和 linger.ms,是保障分区行为可预期的关键。

Kafka Eagle可视化工具
Kafka Eagle可视化工具

Kafka Eagle是一款结合了目前大数据Kafka监控工具的特点,重新研发的一块开源免费的Kafka集群优秀的监控工具。它可以非常方便的监控生产环境中的offset、lag变化、partition分布、owner等,有需要的小伙伴快来保存下载体验吧!

下载

相关标签:

本站声明:本文内容由网友自发贡献,版权归原作者所有,本站不承担相应法律责任。如您发现有涉嫌抄袭侵权的内容,请联系admin@php.cn

相关专题

更多
kafka消费者组有什么作用
kafka消费者组有什么作用

kafka消费者组的作用:1、负载均衡;2、容错性;3、广播模式;4、灵活性;5、自动故障转移和领导者选举;6、动态扩展性;7、顺序保证;8、数据压缩;9、事务性支持。本专题为大家提供相关的文章、下载、课程内容,供大家免费下载体验。

2024.01.12

2286

5

kafka消费组的作用是什么
kafka消费组的作用是什么

kafka消费组的作用:1、负载均衡;2、容错性;3、灵活性;4、高可用性;5、扩展性;6、顺序保证;7、数据压缩;8、事务性支持。本专题为大家提供相关的文章、下载、课程内容,供大家免费下载体验。

2024.02.23

570

5

rabbitmq和kafka有什么区别
rabbitmq和kafka有什么区别

rabbitmq和kafka的区别:1、语言与平台;2、消息传递模型;3、可靠性;4、性能与吞吐量;5、集群与负载均衡;6、消费模型;7、用途与场景;8、社区与生态系统;9、监控与管理;10、其他特性。本专题为大家提供相关的文章、下载、课程内容,供大家免费下载体验。

2024.02.23

544

5

Java 流式处理与 Apache Kafka 实战
Java 流式处理与 Apache Kafka 实战

本专题专注讲解 Java 在流式数据处理与消息队列系统中的应用,系统讲解 Apache Kafka 的基础概念、生产者与消费者模型、Kafka Streams 与 KSQL 流式处理框架、实时数据分析与监控,结合实际业务场景,帮助开发者构建 高吞吐量、低延迟的实时数据流管道,实现高效的数据流转与处理。

2026.02.04

590

32

Buffalo框架数据库开发全教程
Buffalo框架数据库开发全教程

本专题围绕Buffalo框架数据库开发,讲解database.yml多环境配置、soda与fizz迁移生成回滚、模型结构体标签、增删改查与条件查询、一对多与多对多关联、数据校验、回调钩子、事务处理及原生SQL执行能力。

2026.09.23

180

15

Buffalo框架路由与请求处理实操指南
Buffalo框架路由与请求处理实操指南

本专题讲解Buffalo框架路由与请求处理机制,涵盖路由注册与分组、资源路由、Handler编写规范、Context上下文方法、参数绑定、中间件编写挂载、Session与Cookie读写、Flash消息及错误页面定制方法。

2026.09.23

80

15

Buffalo框架零基础入门教程
Buffalo框架零基础入门教程

本专题整理Buffalo框架入门内容,涵盖Go环境准备、buffalo CLI安装、新项目生成、目录结构说明、dev热加载启动、数据库连接配置与常见报错排查,帮助新手按约定优于配置的思路跑通第一个Buffalo框架应用。

2026.09.23

80

15

Conan创建软件包配方指南
Conan创建软件包配方指南

本专题介绍通过conanfile.py创建软件包的方法,讲解包名、版本、依赖和构建设置等基础信息,以及source、build、package、package_info等常用方法的作用及编写思路。

2026.09.22

60

12

Conan二进制包配置指南
Conan二进制包配置指南

本专题介绍Conan根据操作系统、编译器、架构和构建类型生成二进制包的方法,讲解Profile、Settings、Options及Package ID的作用,帮助管理不同平台和编译环境下的包版本。

2026.09.22

60

13

热门下载

更多
网站特效
/
网站源码
/
网站素材
/
前端模板

精品课程

更多
热门推荐
/
最新课程
phpStudy极速入门视频教程
phpStudy极速入门视频教程

共6课时 | 54.6万人学习

独孤九贱(4)_PHP视频教程
独孤九贱(4)_PHP视频教程

共89课时 | 133.3万人学习