You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.04 07:01:07