Kafka 2.4+版本RoundRobinPartitioner分区分配不均问题咨询
Kafka 2.4.1版本RoundRobinPartitioner仅发送消息至偶数分区问题及修复状态
问题现象
在Kafka 2.4.1版本中使用RoundRobinPartitioner时,消息仅被分配到偶数编号的分区;而在Kafka 1.0.0版本中使用相同的分区器逻辑,消息能均匀分配到所有分区。
原因分析
这是Kafka公开BUG KAFKA-9965,由KIP-480的批次优化引发:在启动新批次时,同一条消息的partition方法会被调用两次。由于RoundRobinPartitioner依赖全局递增计数器来选择分区,每次调用都会让计数器加1,最终导致实际分配时跳过了奇数分区。
修复状态
该BUG的相关PR于2021年10月提交后曾停滞,但目前已在Kafka 3.0.0及更高版本中完成修复并正式发布:
- 若使用的是2.4.x系列版本,官方未针对该分支发布补丁,建议升级到3.0.0及以上版本解决问题。
- 若暂时无法升级,可采用以下临时解决方案。
临时解决方案
方案1:自定义幂等RoundRobin分区器
修改分区器逻辑,避免重复递增计数器,例如通过线程局部变量缓存当前消息的分区结果:
public class IdempotentRoundRobinPartitioner implements Partitioner { private final ConcurrentMap<String, AtomicInteger> topicCounterMap = new ConcurrentHashMap<>(); private final ThreadLocal<Integer> currentPartition = new ThreadLocal<>(); @Override public void configure(Map<String, ?> configs) {} @Override public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) { // 检查当前线程是否已有缓存的分区结果 Integer cachedPartition = currentPartition.get(); if (cachedPartition != null) { currentPartition.remove(); return cachedPartition; } List<PartitionInfo> partitions = cluster.partitionsForTopic(topic); int numPartitions = partitions.size(); int nextValue = nextValue(topic); List<PartitionInfo> availablePartitions = cluster.availablePartitionsForTopic(topic); int targetPartition; if (!availablePartitions.isEmpty()) { int part = Utils.toPositive(nextValue) % availablePartitions.size(); targetPartition = availablePartitions.get(part).partition(); } else { targetPartition = Utils.toPositive(nextValue) % numPartitions; } // 缓存分区结果,应对第二次调用 currentPartition.set(targetPartition); return targetPartition; } private int nextValue(String topic) { AtomicInteger counter = topicCounterMap.computeIfAbsent(topic, k -> new AtomicInteger(0)); return counter.getAndIncrement(); } @Override public void close() {} }
方案2:禁用KIP-480批次优化(不推荐)
通过设置batch.size=0关闭批次功能,避免partition方法被重复调用,但此方式会降低Producer性能,仅作为临时应急方案。
代码示例与测试结果
Producer测试代码
package com.example.javakafka; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.serialization.StringSerializer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.util.Properties; public class ProducerDemo { private static final Logger log = LoggerFactory.getLogger(ProducerDemo.class); public static void main(String[] args) { String bootstrapServers = "127.0.0.1:9092"; Properties properties = new Properties(); properties.setProperty(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); properties.setProperty(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); properties.setProperty(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); properties.setProperty("partitioner.class", "com.example.javakafka.RoundRobinPartitionerNew"); KafkaProducer<String, String> producer = new KafkaProducer<>(properties); ProducerRecord<String, String> producerRecord; int counter = 0; while(counter < 40) { producerRecord = new ProducerRecord<>("demo-java-topic-10", "hello world"); producer.send(producerRecord); counter++; } producer.flush(); producer.close(); } }
自定义RoundRobinPartitioner(与官方实现一致)
package com.example.javakafka; import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.atomic.AtomicInteger; import org.apache.kafka.clients.producer.Partitioner; import org.apache.kafka.common.Cluster; import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.utils.Utils; public class RoundRobinPartitionerNew implements Partitioner { private final ConcurrentMap<String, AtomicInteger> topicCounterMap = new ConcurrentHashMap<>(); public RoundRobinPartitionerNew() {} @Override public void configure(Map<String, ?> configs) {} @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(); int nextValue = nextValue(topic); List<PartitionInfo> availablePartitions = cluster.availablePartitionsForTopic(topic); if (!availablePartitions.isEmpty()) { int part = Utils.toPositive(nextValue) % availablePartitions.size(); return availablePartitions.get(part).partition(); } else { return Utils.toPositive(nextValue) % numPartitions; } } private int nextValue(String topic) { AtomicInteger counter = topicCounterMap.computeIfAbsent(topic, k -> new AtomicInteger(0)); return counter.getAndIncrement(); } @Override public void close() {} }
测试结果对比
Kafka 1.0.0版本(正常分配)
CreateTime:1673933856783 Partition:9 null hello world CreateTime:1673933856784 Partition:9 null hello world CreateTime:1673933856784 Partition:9 null hello world CreateTime:1673933856784 Partition:9 null hello world CreateTime:1673933856783 Partition:8 null hello world CreateTime:1673933856784 Partition:8 null hello world CreateTime:1673933856784 Partition:8 null hello world CreateTime:1673933856784 Partition:8 null hello world CreateTime:1673933856783 Partition:5 null hello world CreateTime:1673933856784 Partition:5 null hello world CreateTime:1673933856784 Partition:5 null hello world CreateTime:1673933856784 Partition:5 null hello world // 剩余结果省略,所有分区均有消息分配
Kafka 2.4.0版本(仅偶数分区有消息)
CreateTime:1673933438685 Partition:8 null hello world CreateTime:1673933438688 Partition:8 null hello world CreateTime:1673933438688 Partition:8 null hello world CreateTime:1673933438688 Partition:8 null hello world CreateTime:1673933438688 Partition:4 null hello world CreateTime:1673933438688 Partition:4 null hello world CreateTime:1673933438688 Partition:4 null hello world CreateTime:1673933438688 Partition:4 null hello world // 剩余结果省略,仅偶数编号分区有消息
内容的提问来源于stack exchange,提问作者ctimus
相关产品推荐
相关产品推荐

