This is an automated email from the ASF dual-hosted git repository. jamesshao pushed a commit to branch upsert in repository https://gitbox.apache.org/repos/asf/incubator-pinot.git
commit 70a913cf0821677dbd697dfe8bece675aaae5a84 Author: Bo Zhang <[email protected]> AuthorDate: Mon Oct 28 22:44:29 2019 -0700 Implement FixedPartitionCountPartitioner Reviewers: #streaming_pinot, sjames Reviewed By: #streaming_pinot, sjames Subscribers: sjames Maniphest Tasks: T4265291 Differential Revision: https://code.uberinternal.com/D3516373 --- pinot-grigio/pinot-grigio-coordinator/pom.xml | 5 ++ ...va => FixedPartitionCountBytesPartitioner.java} | 36 +++++------ ...java => FixedPartitionCountIntPartitioner.java} | 30 ++++----- ...er.java => FixedPartitionCountPartitioner.java} | 25 ++++---- .../rpcQueue/KeyCoordinatorQueueProducer.java | 4 +- .../rpcQueue/LogCoordinatorQueueProducer.java | 4 +- .../common/rpcQueue/VersionMsgQueueProducer.java | 4 +- .../keyCoordinator/starter/KeyCoordinatorConf.java | 2 +- .../FixedPartitionCountBytesPartitionerTest.java | 72 +++++++++++++++++++++ .../FixedPartitionCountIntPartitionerTest.java | 74 ++++++++++++++++++++++ .../org.mockito.plugins.MockMaker | 1 + 11 files changed, 203 insertions(+), 54 deletions(-) diff --git a/pinot-grigio/pinot-grigio-coordinator/pom.xml b/pinot-grigio/pinot-grigio-coordinator/pom.xml index 6222a8f..f8f6e48 100644 --- a/pinot-grigio/pinot-grigio-coordinator/pom.xml +++ b/pinot-grigio/pinot-grigio-coordinator/pom.xml @@ -93,6 +93,11 @@ <artifactId>testng</artifactId> <scope>test</scope> </dependency> + <dependency> + <groupId>org.mockito</groupId> + <artifactId>mockito-core</artifactId> + <scope>test</scope> + </dependency> </dependencies> diff --git a/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/IntPartitioner.java b/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/FixedPartitionCountBytesPartitioner.java similarity index 56% copy from pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/IntPartitioner.java copy to pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/FixedPartitionCountBytesPartitioner.java index cd83ee9..9b1de08 100644 --- a/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/IntPartitioner.java +++ b/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/FixedPartitionCountBytesPartitioner.java @@ -18,29 +18,27 @@ */ package org.apache.pinot.grigio.common; -import com.google.common.base.Preconditions; -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; -import java.util.List; -import java.util.Map; -public class IntPartitioner implements Partitioner { - @Override - public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) { - Preconditions.checkState(key instanceof Integer, "expect key to be an integer"); - List<PartitionInfo> partitions = cluster.partitionsForTopic(topic); - int numPartitions = partitions.size(); - return (Integer) key % numPartitions; - } - - @Override - public void close() { - - } +/** + * Fixed partition count partitioner that partition with the bytes of the primary key + */ +public class FixedPartitionCountBytesPartitioner extends FixedPartitionCountPartitioner { @Override - public void configure(Map<String, ?> configs) { + public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) { + if (keyBytes == null) { + throw new IllegalArgumentException("Cannot partition without a key"); + } + int numPartitions = cluster.partitionCountForTopic(topic); + int partitionCount = getPartitionCount(); + if (partitionCount > numPartitions) { + throw new IllegalArgumentException(String + .format("Cannot partition to %d partitions for records in topic %s, which has only %d partitions.", + partitionCount, topic, numPartitions)); + } + return Utils.toPositive(Utils.murmur2(keyBytes)) % partitionCount; } } diff --git a/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/IntPartitioner.java b/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/FixedPartitionCountIntPartitioner.java similarity index 64% copy from pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/IntPartitioner.java copy to pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/FixedPartitionCountIntPartitioner.java index cd83ee9..3969098 100644 --- a/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/IntPartitioner.java +++ b/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/FixedPartitionCountIntPartitioner.java @@ -19,28 +19,24 @@ package org.apache.pinot.grigio.common; import com.google.common.base.Preconditions; -import org.apache.kafka.clients.producer.Partitioner; import org.apache.kafka.common.Cluster; -import org.apache.kafka.common.PartitionInfo; -import java.util.List; -import java.util.Map; -public class IntPartitioner implements Partitioner { +/** + * Fixed partition count partitioner that assumes that the primary key is integer and use it as the partition directly + */ +public class FixedPartitionCountIntPartitioner extends FixedPartitionCountPartitioner { + @Override public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) { Preconditions.checkState(key instanceof Integer, "expect key to be an integer"); - List<PartitionInfo> partitions = cluster.partitionsForTopic(topic); - int numPartitions = partitions.size(); - return (Integer) key % numPartitions; - } - - @Override - public void close() { - - } - - @Override - public void configure(Map<String, ?> configs) { + int numPartitions = cluster.partitionCountForTopic(topic); + int partitionCount = getPartitionCount(); + if (partitionCount > numPartitions) { + throw new IllegalArgumentException(String + .format("Cannot partition to %d partitions for records in topic %s, which has only %d partitions.", + partitionCount, topic, numPartitions)); + } + return (Integer) key % partitionCount; } } diff --git a/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/IntPartitioner.java b/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/FixedPartitionCountPartitioner.java similarity index 65% rename from pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/IntPartitioner.java rename to pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/FixedPartitionCountPartitioner.java index cd83ee9..7e36abb 100644 --- a/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/IntPartitioner.java +++ b/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/FixedPartitionCountPartitioner.java @@ -19,20 +19,21 @@ package org.apache.pinot.grigio.common; import com.google.common.base.Preconditions; +import java.util.Map; import org.apache.kafka.clients.producer.Partitioner; -import org.apache.kafka.common.Cluster; -import org.apache.kafka.common.PartitionInfo; -import java.util.List; -import java.util.Map; -public class IntPartitioner implements Partitioner { - @Override - public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) { - Preconditions.checkState(key instanceof Integer, "expect key to be an integer"); - List<PartitionInfo> partitions = cluster.partitionsForTopic(topic); - int numPartitions = partitions.size(); - return (Integer) key % numPartitions; +/** + * Kafka partitioner that partition records to a fixed number of partitions. + * i.e., partition results will not change even when more partitions are added to the Kafka topic + */ +public abstract class FixedPartitionCountPartitioner implements Partitioner { + + private static final String PARTITION_COUNT = "partition.count"; + private int partitionCount; + + int getPartitionCount() { + return partitionCount; } @Override @@ -42,5 +43,7 @@ public class IntPartitioner implements Partitioner { @Override public void configure(Map<String, ?> configs) { + partitionCount = Integer.parseInt((String)configs.get(PARTITION_COUNT)); + Preconditions.checkState(partitionCount > 0, "Partition count must be greater than 0"); } } diff --git a/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/rpcQueue/KeyCoordinatorQueueProducer.java b/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/rpcQueue/KeyCoordinatorQueueProducer.java index 3cba067..ee417f2 100644 --- a/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/rpcQueue/KeyCoordinatorQueueProducer.java +++ b/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/rpcQueue/KeyCoordinatorQueueProducer.java @@ -22,10 +22,10 @@ import com.google.common.base.Preconditions; import org.apache.commons.configuration.Configuration; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerConfig; -import org.apache.kafka.clients.producer.internals.DefaultPartitioner; import org.apache.kafka.common.serialization.ByteArraySerializer; import org.apache.pinot.grigio.common.CoordinatorConfig; import org.apache.pinot.grigio.common.DistributedCommonUtils; +import org.apache.pinot.grigio.common.FixedPartitionCountBytesPartitioner; import org.apache.pinot.grigio.common.config.CommonConfig; import org.apache.pinot.grigio.common.messages.KeyCoordinatorQueueMsg; import org.apache.pinot.grigio.common.metrics.GrigioMetrics; @@ -67,7 +67,7 @@ public class KeyCoordinatorQueueProducer extends KafkaQueueProducer<byte[], KeyC conf.subset(CoordinatorConfig.KAFKA_CONFIG.KAFKA_CONFIG_KEY)); kafkaProducerConfig.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, ByteArraySerializer.class.getName()); - kafkaProducerConfig.put(ProducerConfig.PARTITIONER_CLASS_CONFIG, DefaultPartitioner.class.getName()); + kafkaProducerConfig.put(ProducerConfig.PARTITIONER_CLASS_CONFIG, FixedPartitionCountBytesPartitioner.class.getName()); DistributedCommonUtils.setKakfaLosslessProducerConfig(kafkaProducerConfig, hostname); _kafkaProducer = new KafkaProducer<>(kafkaProducerConfig); diff --git a/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/rpcQueue/LogCoordinatorQueueProducer.java b/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/rpcQueue/LogCoordinatorQueueProducer.java index 368cf35..9d2b0b7 100644 --- a/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/rpcQueue/LogCoordinatorQueueProducer.java +++ b/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/rpcQueue/LogCoordinatorQueueProducer.java @@ -29,7 +29,7 @@ import org.apache.pinot.grigio.common.metrics.GrigioMetrics; import org.apache.pinot.grigio.common.utils.CommonUtils; import org.apache.pinot.grigio.common.CoordinatorConfig; import org.apache.pinot.grigio.common.DistributedCommonUtils; -import org.apache.pinot.grigio.common.IntPartitioner; +import org.apache.pinot.grigio.common.FixedPartitionCountIntPartitioner; import java.util.Properties; @@ -65,7 +65,7 @@ public class LogCoordinatorQueueProducer extends KafkaQueueProducer<Integer, Log _grigioMetrics = grigioMetrics; kafkaProducerConfig.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, IntegerSerializer.class.getName()); - kafkaProducerConfig.put(ProducerConfig.PARTITIONER_CLASS_CONFIG, IntPartitioner.class.getName()); + kafkaProducerConfig.put(ProducerConfig.PARTITIONER_CLASS_CONFIG, FixedPartitionCountIntPartitioner.class.getName()); DistributedCommonUtils.setKakfaLosslessProducerConfig(kafkaProducerConfig, hostName); this._kafkaProducer = new KafkaProducer<>(kafkaProducerConfig); diff --git a/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/rpcQueue/VersionMsgQueueProducer.java b/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/rpcQueue/VersionMsgQueueProducer.java index 4d1e620..edacd33 100644 --- a/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/rpcQueue/VersionMsgQueueProducer.java +++ b/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/common/rpcQueue/VersionMsgQueueProducer.java @@ -25,7 +25,7 @@ import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.common.serialization.IntegerSerializer; import org.apache.pinot.grigio.common.CoordinatorConfig; import org.apache.pinot.grigio.common.DistributedCommonUtils; -import org.apache.pinot.grigio.common.IntPartitioner; +import org.apache.pinot.grigio.common.FixedPartitionCountIntPartitioner; import org.apache.pinot.grigio.common.config.CommonConfig; import org.apache.pinot.grigio.common.messages.KeyCoordinatorQueueMsg; import org.apache.pinot.grigio.common.metrics.GrigioMetrics; @@ -66,7 +66,7 @@ public class VersionMsgQueueProducer extends KafkaQueueProducer<Integer, KeyCoor conf.subset(CoordinatorConfig.KAFKA_CONFIG.KAFKA_CONFIG_KEY)); kafkaProducerConfig.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, IntegerSerializer.class.getName()); - kafkaProducerConfig.put(ProducerConfig.PARTITIONER_CLASS_CONFIG, IntPartitioner.class.getName()); + kafkaProducerConfig.put(ProducerConfig.PARTITIONER_CLASS_CONFIG, FixedPartitionCountIntPartitioner.class.getName()); DistributedCommonUtils.setKakfaLosslessProducerConfig(kafkaProducerConfig, hostname); _kafkaProducer = new KafkaProducer<>(kafkaProducerConfig); diff --git a/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/keyCoordinator/starter/KeyCoordinatorConf.java b/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/keyCoordinator/starter/KeyCoordinatorConf.java index a48a335..d4ccf02 100644 --- a/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/keyCoordinator/starter/KeyCoordinatorConf.java +++ b/pinot-grigio/pinot-grigio-coordinator/src/main/java/org/apache/pinot/grigio/keyCoordinator/starter/KeyCoordinatorConf.java @@ -62,7 +62,7 @@ public class KeyCoordinatorConf extends PropertiesConfiguration { public static final String KAFKA_CONSUMER_GROUP_ID_PREFIX = "pinot_upsert_kc_consumerGroup_"; private static final String KC_MESSAGE_TOPIC = "kc.message.topic"; - private static final String KC_MESSAGE_PARTITION_COUNT = "kc.message.partition.count"; // todo: get partition count from topic + private static final String KC_MESSAGE_PARTITION_COUNT = "kc.message.partition.count"; public static final String KC_OUTPUT_TOPIC_PREFIX_KEY = "kc.output.topic.prefix"; diff --git a/pinot-grigio/pinot-grigio-coordinator/src/test/java/org/apache/pinot/grigio/common/FixedPartitionCountBytesPartitionerTest.java b/pinot-grigio/pinot-grigio-coordinator/src/test/java/org/apache/pinot/grigio/common/FixedPartitionCountBytesPartitionerTest.java new file mode 100644 index 0000000..9a4b818 --- /dev/null +++ b/pinot-grigio/pinot-grigio-coordinator/src/test/java/org/apache/pinot/grigio/common/FixedPartitionCountBytesPartitionerTest.java @@ -0,0 +1,72 @@ +/** + * 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.pinot.grigio.common; + +import java.util.HashMap; +import java.util.Map; +import org.apache.kafka.common.Cluster; +import org.testng.annotations.BeforeClass; +import org.testng.annotations.Test; + +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; +import static org.testng.Assert.assertEquals; + + +public class FixedPartitionCountBytesPartitionerTest { + + private Map<String, String> configs; + + @BeforeClass + public void setUp() { + configs = new HashMap<>(); + configs.put("partition.count", "4"); + } + + @Test + public void testPartition() { + FixedPartitionCountBytesPartitioner partitioner = new FixedPartitionCountBytesPartitioner(); + partitioner.configure(configs); + + String topic1 = "test-topic1"; + String topic2 = "test-topic2"; + Cluster cluster = mock(Cluster.class); + when(cluster.partitionCountForTopic(topic1)).thenReturn(4); + when(cluster.partitionCountForTopic(topic2)).thenReturn(8); + + String key = "test-key"; + byte[] keyBytes = key.getBytes(); + + int partitionResult1 = partitioner.partition(topic1, key, keyBytes, null, null, cluster); + int partitionResult2 = partitioner.partition(topic2, key, keyBytes, null, null, cluster); + assertEquals(partitionResult1, partitionResult2); + } + + @Test(expectedExceptions = IllegalArgumentException.class) + public void testPartitionFailed() { + FixedPartitionCountBytesPartitioner partitioner = new FixedPartitionCountBytesPartitioner(); + partitioner.configure(configs); + + String topic = "test-topic"; + Cluster cluster = mock(Cluster.class); + when(cluster.partitionCountForTopic(topic)).thenReturn(2); + + partitioner.partition(topic, null, null, null, null, cluster); + } +} diff --git a/pinot-grigio/pinot-grigio-coordinator/src/test/java/org/apache/pinot/grigio/common/FixedPartitionCountIntPartitionerTest.java b/pinot-grigio/pinot-grigio-coordinator/src/test/java/org/apache/pinot/grigio/common/FixedPartitionCountIntPartitionerTest.java new file mode 100644 index 0000000..6581ebb --- /dev/null +++ b/pinot-grigio/pinot-grigio-coordinator/src/test/java/org/apache/pinot/grigio/common/FixedPartitionCountIntPartitionerTest.java @@ -0,0 +1,74 @@ +/** + * 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.pinot.grigio.common; + +import java.util.HashMap; +import java.util.Map; +import org.apache.kafka.common.Cluster; +import org.testng.annotations.BeforeClass; +import org.testng.annotations.Test; + +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; +import static org.testng.Assert.assertEquals; + + +public class FixedPartitionCountIntPartitionerTest { + + private Map<String, String> configs; + + @BeforeClass + public void setUp() { + configs = new HashMap<>(); + configs.put("partition.count", "4"); + } + + @Test + public void testPartition() { + FixedPartitionCountIntPartitioner partitioner = new FixedPartitionCountIntPartitioner(); + partitioner.configure(configs); + + String topic1 = "test-topic1"; + String topic2 = "test-topic2"; + Cluster cluster = mock(Cluster.class); + when(cluster.partitionCountForTopic(topic1)).thenReturn(4); + when(cluster.partitionCountForTopic(topic2)).thenReturn(8); + + Integer key = 6; + + int partitionResult1 = partitioner.partition(topic1, key, null, null, null, cluster); + int partitionResult2 = partitioner.partition(topic2, key, null, null, null, cluster); + assertEquals(partitionResult1, 2); + assertEquals(partitionResult2, 2); + } + + @Test(expectedExceptions = IllegalArgumentException.class) + public void testPartitionFailed() { + FixedPartitionCountIntPartitioner partitioner = new FixedPartitionCountIntPartitioner(); + partitioner.configure(configs); + + String topic = "test-topic"; + Cluster cluster = mock(Cluster.class); + when(cluster.partitionCountForTopic(topic)).thenReturn(2); + + Integer key = 6; + + partitioner.partition(topic, key, null, null, null, cluster); + } +} \ No newline at end of file diff --git a/pinot-grigio/pinot-grigio-coordinator/src/test/resources/mockito-extensions/org.mockito.plugins.MockMaker b/pinot-grigio/pinot-grigio-coordinator/src/test/resources/mockito-extensions/org.mockito.plugins.MockMaker new file mode 100644 index 0000000..ca6ee9c --- /dev/null +++ b/pinot-grigio/pinot-grigio-coordinator/src/test/resources/mockito-extensions/org.mockito.plugins.MockMaker @@ -0,0 +1 @@ +mock-maker-inline \ No newline at end of file --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
