欢迎来到尧图网

客户服务 关于我们

您的位置:首页 > 文旅 > 艺术 > Kafka官方提供的RoundRobinPartitioner出现奇偶数据不均匀

Kafka官方提供的RoundRobinPartitioner出现奇偶数据不均匀

2025/2/21 3:07:16 来源:https://blog.csdn.net/qq_27242695/article/details/139988115  浏览:    关键词:Kafka官方提供的RoundRobinPartitioner出现奇偶数据不均匀

Kafka官方提供的RoundRobinPartitioner出现奇偶数据不均匀

参考:
https://www.cnblogs.com/cbc-onne/p/18140043

  1. 使用RoundRobinPartitioner
/** Licensed to the Apache Software Foundation (ASF) under one or more* contributor license agreements. See the NOTICE file distributed with* this work for additional information regarding copyright ownership.* The ASF licenses this file to You under the Apache License, Version 2.0* (the "License"); you may not use this file except in compliance with* the License. You may obtain a copy of the License at**    http://www.apache.org/licenses/LICENSE-2.0** Unless required by applicable law or agreed to in writing, software* distributed under the License is distributed on an "AS IS" BASIS,* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.* See the License for the specific language governing permissions and* limitations under the License.*/
package org.apache.kafka.clients.producer;import java.util.List;
import java.util.Map;
import java.util.Queue;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentLinkedQueue;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.atomic.AtomicInteger;import org.apache.kafka.common.Cluster;
import org.apache.kafka.common.PartitionInfo;
import org.apache.kafka.common.utils.Utils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;/*** The "Round-Robin" partitioner - MODIFIED TO WORK PROPERLY WITH STICKY PARTITIONING (KIP-480)* <p>* This partitioning strategy can be used when user wants to distribute the writes to all* partitions equally. This is the behaviour regardless of record key hash.*/
public class RoundRobinPartitioner implements Partitioner {private static final Logger LOGGER = LoggerFactory.getLogger(RoundRobinPartitioner.class);private final ConcurrentMap<String, AtomicInteger> topicCounterMap = new ConcurrentHashMap<>();private final ConcurrentMap<String, Queue<Integer>> topicPartitionQueueMap = new ConcurrentHashMap<>();public void configure(Map<String, ?> configs) {}/*** Compute the partition for the given record.** @param topic      The topic name* @param key        The key to partition on (or null if no key)* @param keyBytes   serialized key to partition on (or null if no key)* @param value      The value to partition on or null* @param valueBytes serialized value to partition on or null* @param cluster    The current cluster metadata*/@Overridepublic int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {Queue<Integer> partitionQueue = partitionQueueComputeIfAbsent(topic);Integer queuedPartition = partitionQueue.poll();if (queuedPartition != null) {LOGGER.trace("Partition chosen from queue: {}", queuedPartition);return queuedPartition;} else {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();int partition = availablePartitions.get(part).partition();LOGGER.trace("Partition chosen: {}", partition);return partition;} else {// no partitions are available, give a non-available partitionreturn Utils.toPositive(nextValue) % numPartitions;}}}private int nextValue(String topic) {AtomicInteger counter =topicCounterMap.computeIfAbsent(topic,k -> {return new AtomicInteger(0);});return counter.getAndIncrement();}private Queue<Integer> partitionQueueComputeIfAbsent(String topic) {return topicPartitionQueueMap.computeIfAbsent(topic, k -> {return new ConcurrentLinkedQueue<>();});}public void close() {}/*** Notifies the partitioner a new batch is about to be created. When using the sticky partitioner,* this method can change the chosen sticky partition for the new batch.** @param topic         The topic name* @param cluster       The current cluster metadata* @param prevPartition The partition previously selected for the record that triggered a new*                      batch*/@Overridepublic void onNewBatch(String topic, Cluster cluster, int prevPartition) {LOGGER.trace("New batch so enqueuing partition {} for topic {}", prevPartition, topic);Queue<Integer> partitionQueue = partitionQueueComputeIfAbsent(topic);partitionQueue.add(prevPartition);}
}

版权声明:

本网仅为发布的内容提供存储空间,不对发表、转载的内容提供任何形式的保证。凡本网注明“来源:XXX网络”的作品,均转载自其它媒体,著作权归作者所有,商业转载请联系作者获得授权,非商业转载请注明出处。

我们尊重并感谢每一位作者,均已注明文章来源和作者。如因作品内容、版权或其它问题,请及时与我们联系,联系邮箱:809451989@qq.com,投稿邮箱:809451989@qq.com

热搜词