diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ConsistentHashingStickyKeyConsumerSelector.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ConsistentHashingStickyKeyConsumerSelector.java index f4008ab2e551f..96029d8abe115 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ConsistentHashingStickyKeyConsumerSelector.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ConsistentHashingStickyKeyConsumerSelector.java @@ -22,8 +22,10 @@ import org.apache.pulsar.broker.service.BrokerServiceException.ConsumerAssignException; import org.apache.pulsar.common.util.Murmur3_32Hash; +import java.util.ArrayList; import java.util.Collections; import java.util.Comparator; +import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.NavigableMap; @@ -126,6 +128,25 @@ public Consumer select(byte[] stickyKey) { } } + @Override + public Map> getConsumerKeyHashRanges() { + Map> result = new LinkedHashMap<>(); + rwLock.readLock().lock(); + try { + int start = 0; + for (Map.Entry> entry: hashRing.entrySet()) { + for (Consumer consumer: entry.getValue()) { + result.computeIfAbsent(consumer.consumerName(), key -> new ArrayList<>()) + .add("[" + start + ", " + entry.getKey() + "]"); + } + start = entry.getKey() + 1; + } + } finally { + rwLock.readLock().unlock(); + } + return result; + } + Map> getRangeConsumer() { return Collections.unmodifiableMap(hashRing); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/HashRangeAutoSplitStickyKeyConsumerSelector.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/HashRangeAutoSplitStickyKeyConsumerSelector.java index 18b07f677e275..f74165e7fc9e3 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/HashRangeAutoSplitStickyKeyConsumerSelector.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/HashRangeAutoSplitStickyKeyConsumerSelector.java @@ -21,8 +21,10 @@ import org.apache.pulsar.broker.service.BrokerServiceException.ConsumerAssignException; import org.apache.pulsar.common.util.Murmur3_32Hash; +import java.util.ArrayList; import java.util.Collections; import java.util.HashMap; +import java.util.List; import java.util.Map; import java.util.Map.Entry; import java.util.concurrent.ConcurrentSkipListMap; @@ -112,6 +114,18 @@ public Consumer select(byte[] stickyKey) { } } + @Override + public Map> getConsumerKeyHashRanges() { + Map> result = new HashMap<>(); + int start = 0; + for (Map.Entry entry: rangeMap.entrySet()) { + result.computeIfAbsent(entry.getValue().consumerName(), key -> new ArrayList<>()) + .add("[" + start + ", " + entry.getKey() + "]"); + start = entry.getKey() + 1; + } + return result; + } + private int findBiggestRange() { int slots = 0; int busiestRange = rangeSize; @@ -147,10 +161,6 @@ private boolean is2Power(int num) { return (num & num - 1) == 0; } - Map getConsumerRange() { - return Collections.unmodifiableMap(consumerRange); - } - Map getRangeConsumer() { return Collections.unmodifiableMap(rangeMap); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/HashRangeExclusiveStickyKeyConsumerSelector.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/HashRangeExclusiveStickyKeyConsumerSelector.java index 0276f42214951..75cb6b5b5adf8 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/HashRangeExclusiveStickyKeyConsumerSelector.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/HashRangeExclusiveStickyKeyConsumerSelector.java @@ -21,11 +21,18 @@ import org.apache.pulsar.common.api.proto.PulsarApi; import org.apache.pulsar.common.util.Murmur3_32Hash; +import java.util.ArrayList; import java.util.Collections; +import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentSkipListMap; +/** + * This is a sticky-key consumer selector based user provided range. + * User is responsible for making sure provided range for all consumers cover the rangeSize + * else there'll be chance that a key fall in a `whole` that not handled by any consumer. + */ public class HashRangeExclusiveStickyKeyConsumerSelector implements StickyKeyConsumerSelector { private final int rangeSize; @@ -63,6 +70,24 @@ public Consumer select(byte[] stickyKey) { return select(Murmur3_32Hash.getInstance().makeHash(stickyKey)); } + @Override + public Map> getConsumerKeyHashRanges() { + Map> result = new HashMap<>(); + Map.Entry prev = null; + for (Map.Entry entry: rangeMap.entrySet()) { + if (prev == null) { + prev = entry; + } else { + if (prev.getValue().equals(entry.getValue())) { + result.computeIfAbsent(entry.getValue().consumerName(), key -> new ArrayList<>()) + .add("[" + prev.getKey() + ", " + entry.getKey() + "]"); + } + prev = null; + } + } + return result; + } + Consumer select(int hash) { if (rangeMap.size() > 0) { int slot = hash % rangeSize; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/StickyKeyConsumerSelector.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/StickyKeyConsumerSelector.java index 1b168d5128f5f..0d686f716faaa 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/StickyKeyConsumerSelector.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/StickyKeyConsumerSelector.java @@ -20,6 +20,9 @@ import org.apache.pulsar.broker.service.BrokerServiceException.ConsumerAssignException; +import java.util.List; +import java.util.Map; + public interface StickyKeyConsumerSelector { int DEFAULT_RANGE_SIZE = 2 << 15; @@ -43,4 +46,10 @@ public interface StickyKeyConsumerSelector { * @return consumer */ Consumer select(byte[] stickyKey); + + /** + * Get key hash ranges handled by each consumer + * @return A map where key is a consumer name and value is list of hash range it receiving message for. + */ + Map> getConsumerKeyHashRanges(); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java index 56963e8a85a91..c8592ba0f4620 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStickyKeyDispatcherMultipleConsumers.java @@ -369,6 +369,10 @@ public LinkedHashMap getRecentlyJoinedConsumers() { return recentlyJoinedConsumers; } + public Map> getConsumerKeyHashRanges() { + return selector.getConsumerKeyHashRanges(); + } + private static final Logger log = LoggerFactory.getLogger(PersistentStickyKeyDispatcherMultipleConsumers.class); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java index 00c8a587938a7..dc78c03d7fd3e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java @@ -480,6 +480,8 @@ public String getTypeString() { return "Failover"; case Shared: return "Shared"; + case Key_Shared: + return "Key_Shared"; } return "Null"; @@ -899,6 +901,8 @@ public SubscriptionStats getStats(Boolean getPreciseBacklog) { subStats.lastConsumedFlowTimestamp = lastConsumedFlowTimestamp; Dispatcher dispatcher = this.dispatcher; if (dispatcher != null) { + Map> consumerKeyHashRanges = getType() == SubType.Key_Shared? + ((PersistentStickyKeyDispatcherMultipleConsumers)dispatcher).getConsumerKeyHashRanges(): null; dispatcher.getConsumers().forEach(consumer -> { ConsumerStats consumerStats = consumer.getStats(); subStats.consumers.add(consumerStats); @@ -911,6 +915,9 @@ public SubscriptionStats getStats(Boolean getPreciseBacklog) { subStats.unackedMessages += consumerStats.unackedMessages; subStats.lastConsumedTimestamp = Math.max(subStats.lastConsumedTimestamp, consumerStats.lastConsumedTimestamp); subStats.lastAckedTimestamp = Math.max(subStats.lastAckedTimestamp, consumerStats.lastAckedTimestamp); + if (consumerKeyHashRanges != null && consumerKeyHashRanges.containsKey(consumer.consumerName())) { + consumerStats.keyHashRanges = consumerKeyHashRanges.get(consumer.consumerName()); + } }); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ConsistentHashingStickyKeyConsumerSelectorTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ConsistentHashingStickyKeyConsumerSelectorTest.java index 53df067f0ff7a..c2b36fac4718a 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ConsistentHashingStickyKeyConsumerSelectorTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ConsistentHashingStickyKeyConsumerSelectorTest.java @@ -21,10 +21,15 @@ import static org.mockito.Mockito.mock; import static org.mockito.Mockito.when; +import java.util.Arrays; import java.util.HashMap; +import java.util.List; import java.util.Map; +import java.util.Set; import java.util.UUID; +import java.util.stream.Collectors; +import com.google.common.collect.ImmutableSet; import org.apache.pulsar.broker.service.BrokerServiceException.ConsumerAssignException; import org.testng.Assert; import org.testng.annotations.Test; @@ -134,4 +139,28 @@ public void testConsumerSelect() throws ConsumerAssignException { Assert.assertEquals(selectionMap.get("c4").intValue(), N); } + + @Test + public void testGetConsumerKeyHashRanges() throws BrokerServiceException.ConsumerAssignException { + ConsistentHashingStickyKeyConsumerSelector selector = new ConsistentHashingStickyKeyConsumerSelector(3); + List consumerName = Arrays.asList("consumer1", "consumer2", "consumer3"); + List range = Arrays.asList(new int[] {0, 2}, new int[] {3, 7}, new int[] {9, 12}, new int[] {15, 20}); + for (int index = 0; index < consumerName.size(); index++) { + Consumer consumer = mock(Consumer.class); + when(consumer.consumerName()).thenReturn(consumerName.get(index)); + selector.addConsumer(consumer); + } + + Map> expectedResult = new HashMap<>(); + expectedResult.put("consumer1", ImmutableSet.of("[0, 330121749]", "[330121750, 618146114]", "[1797637922, 1976098885]")); + expectedResult.put("consumer2", ImmutableSet.of("[938427576, 1094135919]", "[1138613629, 1342907082]", "[1342907083, 1797637921]")); + expectedResult.put("consumer3", ImmutableSet.of("[618146115, 772640562]", "[772640563, 938427575]", "[1094135920, 1138613628]")); + for (Map.Entry> entry : selector.getConsumerKeyHashRanges().entrySet()) { + Assert.assertEquals(entry.getValue().stream().collect(Collectors.toSet()), expectedResult.get(entry.getKey())); + expectedResult.remove(entry.getKey()); + } + Assert.assertEquals(expectedResult.size(), 0); + } + + } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/HashRangeAutoSplitStickyKeyConsumerSelectorTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/HashRangeAutoSplitStickyKeyConsumerSelectorTest.java new file mode 100644 index 0000000000000..450b92b9fa307 --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/HashRangeAutoSplitStickyKeyConsumerSelectorTest.java @@ -0,0 +1,58 @@ +/** + * 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.pulsar.broker.service; + +import com.google.common.collect.ImmutableList; +import org.apache.pulsar.common.api.proto.PulsarApi; +import org.testng.Assert; +import org.testng.annotations.Test; + +import java.util.Arrays; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +public class HashRangeAutoSplitStickyKeyConsumerSelectorTest { + + @Test + public void testGetConsumerKeyHashRanges() throws BrokerServiceException.ConsumerAssignException { + HashRangeAutoSplitStickyKeyConsumerSelector selector = new HashRangeAutoSplitStickyKeyConsumerSelector(2 << 5); + List consumerName = Arrays.asList("consumer1", "consumer2", "consumer3", "consumer4"); + for (int index = 0; index < consumerName.size(); index++) { + Consumer consumer = mock(Consumer.class); + when(consumer.consumerName()).thenReturn(consumerName.get(index)); + selector.addConsumer(consumer); + } + + Map> expectedResult = new HashMap<>(); + expectedResult.put("consumer1", ImmutableList.of("[49, 64]")); + expectedResult.put("consumer4", ImmutableList.of("[33, 48]")); + expectedResult.put("consumer2", ImmutableList.of("[17, 32]")); + expectedResult.put("consumer3", ImmutableList.of("[0, 16]")); + for (Map.Entry> entry : selector.getConsumerKeyHashRanges().entrySet()) { + Assert.assertEquals(entry.getValue(), expectedResult.get(entry.getKey())); + expectedResult.remove(entry.getKey()); + } + Assert.assertEquals(expectedResult.size(), 0); + } + +} diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/HashRangeExclusiveStickyKeyConsumerSelectorTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/HashRangeExclusiveStickyKeyConsumerSelectorTest.java index 97d24e7ff8d53..816d0133a088c 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/HashRangeExclusiveStickyKeyConsumerSelectorTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/HashRangeExclusiveStickyKeyConsumerSelectorTest.java @@ -18,13 +18,17 @@ */ package org.apache.pulsar.broker.service; +import com.google.common.collect.ImmutableList; import com.google.common.collect.Lists; import org.apache.pulsar.common.api.proto.PulsarApi; import org.testng.Assert; import org.testng.annotations.Test; import java.util.ArrayList; +import java.util.Arrays; +import java.util.HashMap; import java.util.List; +import java.util.Map; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.when; @@ -107,6 +111,36 @@ public void testInvalidRangeTotal() { new HashRangeExclusiveStickyKeyConsumerSelector(0); } + @Test + public void testGetConsumerKeyHashRanges() throws BrokerServiceException.ConsumerAssignException { + HashRangeExclusiveStickyKeyConsumerSelector selector = new HashRangeExclusiveStickyKeyConsumerSelector(10); + List consumerName = Arrays.asList("consumer1", "consumer2", "consumer3", "consumer4"); + List range = Arrays.asList(new int[] {0, 2}, new int[] {3, 7}, new int[] {9, 12}, new int[] {15, 20}); + for (int index = 0; index < consumerName.size(); index++) { + Consumer consumer = mock(Consumer.class); + PulsarApi.KeySharedMeta keySharedMeta = PulsarApi.KeySharedMeta.newBuilder() + .setKeySharedMode(PulsarApi.KeySharedMode.STICKY) + .addHashRanges(PulsarApi.IntRange.newBuilder().setStart(range.get(index)[0]) + .setEnd(range.get(index)[1]).build()) + .build(); + when(consumer.getKeySharedMeta()).thenReturn(keySharedMeta); + when(consumer.consumerName()).thenReturn(consumerName.get(index)); + Assert.assertEquals(consumer.getKeySharedMeta(), keySharedMeta); + selector.addConsumer(consumer); + } + + Map> expectedResult = new HashMap<>(); + expectedResult.put("consumer1", ImmutableList.of("[0, 2]")); + expectedResult.put("consumer2", ImmutableList.of("[3, 7]")); + expectedResult.put("consumer3", ImmutableList.of("[9, 12]")); + expectedResult.put("consumer4", ImmutableList.of("[15, 20]")); + for (Map.Entry> entry : selector.getConsumerKeyHashRanges().entrySet()) { + Assert.assertEquals(entry.getValue(), expectedResult.get(entry.getKey())); + expectedResult.remove(entry.getKey()); + } + Assert.assertEquals(expectedResult.size(), 0); + } + @Test public void testSingleRangeConflict() throws BrokerServiceException.ConsumerAssignException { HashRangeExclusiveStickyKeyConsumerSelector selector = new HashRangeExclusiveStickyKeyConsumerSelector(10); diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/ConsumerStats.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/ConsumerStats.java index 837390f73158c..f1836854a8983 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/ConsumerStats.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/ConsumerStats.java @@ -20,6 +20,7 @@ import static com.google.common.base.Preconditions.checkNotNull; +import java.util.List; import java.util.Map; /** @@ -77,6 +78,9 @@ public class ConsumerStats { public long lastAckedTimestamp; public long lastConsumedTimestamp; + /** Hash ranges assigned to this consumer if is Key_Shared sub mode. **/ + public List keyHashRanges; + /** Metadata (key/value strings) associated with this consumer. */ public Map metadata;