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


















