讲师中心 微信公众号
AI工具推荐 视频效率加速

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

云浩吖_7732

云浩吖_7732

发布时间:2026-09-29 16:25:40

|

125人浏览过

|

来源于php中文网

原创

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;
}

并在 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);
}

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

热门AI工具

更多
二狗PPT
二狗PPT Hot

一款AI演示文稿工具,主要用于专为中式职场打造的AI PPT生成工具,适合需要提升相关任务效率的用户。

Loomy
Loomy Hot

一款AI工具,主要用于科大讯飞发布的桌面级 AI 助理,比 OpenClaw 更易用、更安全!,适合需要提升相关任务效率的用户。

Laper
Laper Hot

Laper是专为编剧、导演和制片人推出的 AI 原生剧本创作工具。

立刻MV
立刻MV Hot

立刻MV是一款AI文本写作工具,AI 音乐视频(MV)创作工具。

切问学术

切问学术是一款AI论文写作工具,复旦大学NLP团队推出的AI学术智能体。

DeepSeek

DeepSeek是一款面向对话、写作、编程和推理场景的AI大模型工具。

豆包大模型

豆包大模型是一款由字节跳动推出的企业级大语言模型服务平台。

蛙蛙写作

一款AI论文写作工具,主要用于超级AI智能写作助手,适合需要提升相关任务效率的用户。

WorkBuddy

一款AI办公效率工具,主要用于腾讯云推出的AI原生桌面智能体工作台,适合需要提升相关任务效率的用户。

相关专题

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

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

2306

2024.01.12

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

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

570

2024.02.23

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

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

544

2024.02.23

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

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

590

2026.02.04

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

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

180

2026.09.23

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

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

80

2026.09.23

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

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

80

2026.09.23

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

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

60

2026.09.22

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

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

60

2026.09.22

热门下载

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

精品课程

更多
热门推荐
/
最新课程
关于我们 免责申明 举报中心 意见反馈 讲师合作 广告合作 最新更新
php中文网:公益在线php培训,帮助PHP学习者快速成长!
关注服务号
PHP中文网订阅号
每天精选资源文章推送

Copyright 2014-2026 https://www.php.cn/ All Rights Reserved | php.cn