From 7e91ea128048c87475117705059452b327b24987 Mon Sep 17 00:00:00 2001 From: Marvin Cai Date: Sat, 14 Nov 2020 21:25:48 -0800 Subject: [PATCH 1/5] Add key-shared comsuner range to internal topic stats. --- ...stentHashingStickyKeyConsumerSelector.java | 19 ++++++++++ ...ngeAutoSplitStickyKeyConsumerSelector.java | 16 ++++++--- ...ngeExclusiveStickyKeyConsumerSelector.java | 23 ++++++++++++ .../service/StickyKeyConsumerSelector.java | 8 +++++ ...tStickyKeyDispatcherMultipleConsumers.java | 4 +++ .../persistent/PersistentSubscription.java | 2 ++ .../service/persistent/PersistentTopic.java | 4 +++ ...tHashingStickyKeyConsumerSelectorTest.java | 29 +++++++++++++++ ...utoSplitStickyKeyConsumerSelectorTest.java | 36 +++++++++++++++++++ ...xclusiveStickyKeyConsumerSelectorTest.java | 28 +++++++++++++++ .../data/PersistentTopicInternalStats.java | 2 ++ 11 files changed, 167 insertions(+), 4 deletions(-) create mode 100644 pulsar-broker/src/test/java/org/apache/pulsar/broker/service/HashRangeAutoSplitStickyKeyConsumerSelectorTest.java 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..19db5c8a6a7dc 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 @@ -24,12 +24,14 @@ import java.util.Collections; import java.util.Comparator; +import java.util.LinkedHashMap; import java.util.List; import java.util.Map; import java.util.NavigableMap; import java.util.TreeMap; import java.util.concurrent.locks.ReadWriteLock; import java.util.concurrent.locks.ReentrantReadWriteLock; +import java.util.stream.Collectors; /** * This is a consumer selector based fixed hash range. @@ -126,6 +128,23 @@ public Consumer select(byte[] stickyKey) { } } + @Override + public Map getConsumerRange() { + Map result = new LinkedHashMap<>(); + rwLock.readLock().lock(); + try { + int start = 0; + for (Map.Entry> entry: hashRing.entrySet()) { + result.put(start + "--" + entry.getKey(), entry.getValue().stream().map( + consumer -> consumer.consumerName()).collect(Collectors.joining(", "))); + 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..d0f700ed54e15 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 @@ -23,6 +23,7 @@ import java.util.Collections; import java.util.HashMap; +import java.util.LinkedHashMap; import java.util.Map; import java.util.Map.Entry; import java.util.concurrent.ConcurrentSkipListMap; @@ -112,6 +113,17 @@ public Consumer select(byte[] stickyKey) { } } + @Override + public Map getConsumerRange() { + Map result = new LinkedHashMap<>(); + int start = 0; + for (Map.Entry entry: rangeMap.entrySet()) { + result.put(start + "--" + entry.getKey(), entry.getValue().consumerName()); + start = entry.getKey() + 1; + } + return result; + } + private int findBiggestRange() { int slots = 0; int busiestRange = rangeSize; @@ -147,10 +159,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..34c48f702ab6a 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 @@ -22,10 +22,16 @@ import org.apache.pulsar.common.util.Murmur3_32Hash; import java.util.Collections; +import java.util.LinkedHashMap; 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 +69,23 @@ public Consumer select(byte[] stickyKey) { return select(Murmur3_32Hash.getInstance().makeHash(stickyKey)); } + @Override + public Map getConsumerRange() { + Map result = new LinkedHashMap<>(); + Map.Entry prev = null; + for (Map.Entry entry: rangeMap.entrySet()) { + if (prev == null) { + prev = entry; + } else { + if (prev.getValue().equals(entry.getValue())) { + result.put(prev.getKey() + "--" + entry.getKey(), entry.getValue().consumerName()); + } + 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..4ab5b7beeadab 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,8 @@ import org.apache.pulsar.broker.service.BrokerServiceException.ConsumerAssignException; +import java.util.Map; + public interface StickyKeyConsumerSelector { int DEFAULT_RANGE_SIZE = 2 << 15; @@ -43,4 +45,10 @@ public interface StickyKeyConsumerSelector { * @return consumer */ Consumer select(byte[] stickyKey); + + /** + * Get range handled by each consumer + * @return A map where key is a range and value is consumer receiving message for the range. + */ + Map getConsumerRange(); } 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..933a368165145 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 getConsumerRange() { + return selector.getConsumerRange(); + } + 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..60c627526503a 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"; diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index ec7215e56d3a0..9a939160db60a 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -1706,6 +1706,10 @@ public CompletableFuture getInternalStats(boolean cs.numberOfEntriesSinceFirstNotAckedMessage = cursor.getNumberOfEntriesSinceFirstNotAckedMessage(); cs.totalNonContiguousDeletedMessagesRange = cursor.getTotalNonContiguousDeletedMessagesRange(); cs.properties = cursor.getProperties(); + Subscription sub = subscriptions.get(Codec.decode(cursor.getName())); + if (sub.getType() == SubType.Key_Shared) { + cs.consumerRange = ((PersistentStickyKeyDispatcherMultipleConsumers)sub.getDispatcher()).getConsumerRange(); + } stats.cursors.put(cursor.getName(), cs); }); if (futures != null) { 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..ab3b73a4fa110 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,11 +21,14 @@ 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.UUID; import org.apache.pulsar.broker.service.BrokerServiceException.ConsumerAssignException; +import org.apache.pulsar.common.api.proto.PulsarApi; import org.testng.Assert; import org.testng.annotations.Test; @@ -134,4 +137,30 @@ public void testConsumerSelect() throws ConsumerAssignException { Assert.assertEquals(selectionMap.get("c4").intValue(), N); } + + @Test + public void testGetConsumerRange() 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); + } + + int index = 0; + List expectedConsumerName = Arrays.asList("consumer1", "consumer1", "consumer3", "consumer3", "consumer2" + , "consumer3", "consumer2", "consumer2", "consumer1"); + List expectedRange = Arrays.asList(new int[] {0, 330121749}, new int[] {330121750, 618146114}, new int[] {618146115, 772640562}, + new int[] {772640563, 938427575}, new int[] {938427576, 1094135919}, new int[] {1094135920, 1138613628}, new int[] {1138613629, 1342907082}, + new int[] {1342907083, 1797637921}, new int[] {1797637922, 1976098885}); + for (Map.Entry entry : selector.getConsumerRange().entrySet()) { + Assert.assertEquals(entry.getKey(), expectedRange.get(index)[0] + "--" + expectedRange.get(index)[1]); + Assert.assertEquals(entry.getValue(), expectedConsumerName.get(index)); + index++; + } + } + + } 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..b7dccc5dce46a --- /dev/null +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/HashRangeAutoSplitStickyKeyConsumerSelectorTest.java @@ -0,0 +1,36 @@ +package org.apache.pulsar.broker.service; + +import org.apache.pulsar.common.api.proto.PulsarApi; +import org.testng.Assert; +import org.testng.annotations.Test; + +import java.util.Arrays; +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 testGetConsumerRange() 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); + } + + int index = 0; + List expectedConsumerName = Arrays.asList("consumer3", "consumer2", "consumer4", "consumer1"); + List expectedRange = Arrays.asList(new int[] {0, 16}, new int[] {17, 32}, new int[] {33, 48}, new int[] {49, 64}); + for (Map.Entry entry : selector.getConsumerRange().entrySet()) { + Assert.assertEquals(entry.getKey(), expectedRange.get(index)[0] + "--" + expectedRange.get(index)[1]); + Assert.assertEquals(entry.getValue(), expectedConsumerName.get(index)); + index++; + } + } + +} 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..ab720e064df12 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 @@ -24,7 +24,9 @@ import org.testng.annotations.Test; import java.util.ArrayList; +import java.util.Arrays; import java.util.List; +import java.util.Map; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.when; @@ -107,6 +109,32 @@ public void testInvalidRangeTotal() { new HashRangeExclusiveStickyKeyConsumerSelector(0); } + @Test + public void testGetConsumerRange() 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); + } + + int index = 0; + for (Map.Entry entry : selector.getConsumerRange().entrySet()) { + Assert.assertEquals(entry.getKey(), range.get(index)[0] + "--" + range.get(index)[1]); + Assert.assertEquals(entry.getValue(), consumerName.get(index)); + index++; + } + } + @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/PersistentTopicInternalStats.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/PersistentTopicInternalStats.java index e5279803e992e..bc0d55b1af81b 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/PersistentTopicInternalStats.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/PersistentTopicInternalStats.java @@ -76,6 +76,8 @@ public static class CursorStats { public long numberOfEntriesSinceFirstNotAckedMessage; public int totalNonContiguousDeletedMessagesRange; + public Map consumerRange; + public Map properties; } } From 5b2fbd19b69a4e414d65ba6b5f3552d13d90fdea Mon Sep 17 00:00:00 2001 From: Marvin Cai Date: Sat, 14 Nov 2020 21:42:03 -0800 Subject: [PATCH 2/5] Add missing license. --- ...AutoSplitStickyKeyConsumerSelectorTest.java | 18 ++++++++++++++++++ 1 file changed, 18 insertions(+) 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 index b7dccc5dce46a..85a77724a9edc 100644 --- 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 @@ -1,3 +1,21 @@ +/** + * 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 org.apache.pulsar.common.api.proto.PulsarApi; From 56dbafece10fb76cd56e906de0f6cabd2b82a1ce Mon Sep 17 00:00:00 2001 From: Marvin Cai Date: Sun, 15 Nov 2020 22:45:52 -0800 Subject: [PATCH 3/5] Move consumer range info to ConsumerStats and expose via topic-stats. --- ...stentHashingStickyKeyConsumerSelector.java | 12 +++++----- ...ngeAutoSplitStickyKeyConsumerSelector.java | 10 +++++---- ...ngeExclusiveStickyKeyConsumerSelector.java | 10 +++++---- .../service/StickyKeyConsumerSelector.java | 5 +++-- ...tStickyKeyDispatcherMultipleConsumers.java | 2 +- .../persistent/PersistentSubscription.java | 5 +++++ .../service/persistent/PersistentTopic.java | 4 ---- ...tHashingStickyKeyConsumerSelectorTest.java | 22 +++++++++---------- ...utoSplitStickyKeyConsumerSelectorTest.java | 18 +++++++++------ ...xclusiveStickyKeyConsumerSelectorTest.java | 16 +++++++++----- .../common/policies/data/ConsumerStats.java | 4 ++++ .../data/PersistentTopicInternalStats.java | 2 -- 12 files changed, 65 insertions(+), 45 deletions(-) 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 19db5c8a6a7dc..37049aab1725e 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,6 +22,7 @@ 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; @@ -31,7 +32,6 @@ import java.util.TreeMap; import java.util.concurrent.locks.ReadWriteLock; import java.util.concurrent.locks.ReentrantReadWriteLock; -import java.util.stream.Collectors; /** * This is a consumer selector based fixed hash range. @@ -129,14 +129,16 @@ public Consumer select(byte[] stickyKey) { } @Override - public Map getConsumerRange() { - Map result = new LinkedHashMap<>(); + public Map> getConsumerRange() { + Map> result = new LinkedHashMap<>(); rwLock.readLock().lock(); try { int start = 0; for (Map.Entry> entry: hashRing.entrySet()) { - result.put(start + "--" + entry.getKey(), entry.getValue().stream().map( - consumer -> consumer.consumerName()).collect(Collectors.joining(", "))); + for (Consumer consumer: entry.getValue()) { + result.computeIfAbsent(consumer.consumerName(), key -> new ArrayList<>()) + .add(start + "--" + entry.getKey()); + } start = entry.getKey() + 1; } } finally { 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 d0f700ed54e15..23de1ff4dd2b2 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,9 +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.LinkedHashMap; +import java.util.List; import java.util.Map; import java.util.Map.Entry; import java.util.concurrent.ConcurrentSkipListMap; @@ -114,11 +115,12 @@ public Consumer select(byte[] stickyKey) { } @Override - public Map getConsumerRange() { - Map result = new LinkedHashMap<>(); + public Map> getConsumerRange() { + Map> result = new HashMap<>(); int start = 0; for (Map.Entry entry: rangeMap.entrySet()) { - result.put(start + "--" + entry.getKey(), entry.getValue().consumerName()); + result.computeIfAbsent(entry.getValue().consumerName(), key -> new ArrayList<>()) + .add(start + "--" + entry.getKey()); start = entry.getKey() + 1; } return result; 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 34c48f702ab6a..91c260c365d6d 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,8 +21,9 @@ 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.LinkedHashMap; +import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentSkipListMap; @@ -70,15 +71,16 @@ public Consumer select(byte[] stickyKey) { } @Override - public Map getConsumerRange() { - Map result = new LinkedHashMap<>(); + public Map> getConsumerRange() { + 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.put(prev.getKey() + "--" + entry.getKey(), entry.getValue().consumerName()); + result.computeIfAbsent(entry.getValue().consumerName(), key -> new ArrayList<>()) + .add(prev.getKey() + "--" + entry.getKey()); } prev = null; } 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 4ab5b7beeadab..3c7f73ebd23b1 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,7 @@ import org.apache.pulsar.broker.service.BrokerServiceException.ConsumerAssignException; +import java.util.List; import java.util.Map; public interface StickyKeyConsumerSelector { @@ -48,7 +49,7 @@ public interface StickyKeyConsumerSelector { /** * Get range handled by each consumer - * @return A map where key is a range and value is consumer receiving message for the range. + * @return A map where key is a consumer name and value is list of range it receiving message for. */ - Map getConsumerRange(); + Map> getConsumerRange(); } 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 933a368165145..8dac8889955bc 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,7 +369,7 @@ public LinkedHashMap getRecentlyJoinedConsumers() { return recentlyJoinedConsumers; } - public Map getConsumerRange() { + public Map> getConsumerRange() { return selector.getConsumerRange(); } 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 60c627526503a..7d8e40bec4b50 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 @@ -901,6 +901,8 @@ public SubscriptionStats getStats(Boolean getPreciseBacklog) { subStats.lastConsumedFlowTimestamp = lastConsumedFlowTimestamp; Dispatcher dispatcher = this.dispatcher; if (dispatcher != null) { + Map> consumerRanges = getType() == SubType.Key_Shared? + ((PersistentStickyKeyDispatcherMultipleConsumers)dispatcher).getConsumerRange(): null; dispatcher.getConsumers().forEach(consumer -> { ConsumerStats consumerStats = consumer.getStats(); subStats.consumers.add(consumerStats); @@ -913,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 (consumerRanges != null && consumerRanges.containsKey(consumer.consumerName())) { + consumerStats.keyHashRange = consumerRanges.get(consumer.consumerName()); + } }); } diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 9a939160db60a..ec7215e56d3a0 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -1706,10 +1706,6 @@ public CompletableFuture getInternalStats(boolean cs.numberOfEntriesSinceFirstNotAckedMessage = cursor.getNumberOfEntriesSinceFirstNotAckedMessage(); cs.totalNonContiguousDeletedMessagesRange = cursor.getTotalNonContiguousDeletedMessagesRange(); cs.properties = cursor.getProperties(); - Subscription sub = subscriptions.get(Codec.decode(cursor.getName())); - if (sub.getType() == SubType.Key_Shared) { - cs.consumerRange = ((PersistentStickyKeyDispatcherMultipleConsumers)sub.getDispatcher()).getConsumerRange(); - } stats.cursors.put(cursor.getName(), cs); }); if (futures != null) { 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 ab3b73a4fa110..7741cf7dfbe0d 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 @@ -25,10 +25,12 @@ 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.apache.pulsar.common.api.proto.PulsarApi; import org.testng.Assert; import org.testng.annotations.Test; @@ -149,17 +151,15 @@ public void testGetConsumerRange() throws BrokerServiceException.ConsumerAssignE selector.addConsumer(consumer); } - int index = 0; - List expectedConsumerName = Arrays.asList("consumer1", "consumer1", "consumer3", "consumer3", "consumer2" - , "consumer3", "consumer2", "consumer2", "consumer1"); - List expectedRange = Arrays.asList(new int[] {0, 330121749}, new int[] {330121750, 618146114}, new int[] {618146115, 772640562}, - new int[] {772640563, 938427575}, new int[] {938427576, 1094135919}, new int[] {1094135920, 1138613628}, new int[] {1138613629, 1342907082}, - new int[] {1342907083, 1797637921}, new int[] {1797637922, 1976098885}); - for (Map.Entry entry : selector.getConsumerRange().entrySet()) { - Assert.assertEquals(entry.getKey(), expectedRange.get(index)[0] + "--" + expectedRange.get(index)[1]); - Assert.assertEquals(entry.getValue(), expectedConsumerName.get(index)); - index++; + 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.getConsumerRange().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 index 85a77724a9edc..6bdab4fb9fd8f 100644 --- 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 @@ -18,11 +18,13 @@ */ 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; @@ -41,14 +43,16 @@ public void testGetConsumerRange() throws BrokerServiceException.ConsumerAssignE selector.addConsumer(consumer); } - int index = 0; - List expectedConsumerName = Arrays.asList("consumer3", "consumer2", "consumer4", "consumer1"); - List expectedRange = Arrays.asList(new int[] {0, 16}, new int[] {17, 32}, new int[] {33, 48}, new int[] {49, 64}); - for (Map.Entry entry : selector.getConsumerRange().entrySet()) { - Assert.assertEquals(entry.getKey(), expectedRange.get(index)[0] + "--" + expectedRange.get(index)[1]); - Assert.assertEquals(entry.getValue(), expectedConsumerName.get(index)); - index++; + 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.getConsumerRange().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 ab720e064df12..f602d294ce9ac 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,6 +18,7 @@ */ 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; @@ -25,6 +26,7 @@ import java.util.ArrayList; import java.util.Arrays; +import java.util.HashMap; import java.util.List; import java.util.Map; @@ -127,12 +129,16 @@ public void testGetConsumerRange() throws BrokerServiceException.ConsumerAssignE selector.addConsumer(consumer); } - int index = 0; - for (Map.Entry entry : selector.getConsumerRange().entrySet()) { - Assert.assertEquals(entry.getKey(), range.get(index)[0] + "--" + range.get(index)[1]); - Assert.assertEquals(entry.getValue(), consumerName.get(index)); - index++; + 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.getConsumerRange().entrySet()) { + Assert.assertEquals(entry.getValue(), expectedResult.get(entry.getKey())); + expectedResult.remove(entry.getKey()); } + Assert.assertEquals(expectedResult.size(), 0); } @Test 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..a7673d0eadc79 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 range assigned to this consumer if is Key_Shared sub mode. **/ + public List keyHashRange; + /** Metadata (key/value strings) associated with this consumer. */ public Map metadata; diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/PersistentTopicInternalStats.java b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/PersistentTopicInternalStats.java index bc0d55b1af81b..e5279803e992e 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/PersistentTopicInternalStats.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/policies/data/PersistentTopicInternalStats.java @@ -76,8 +76,6 @@ public static class CursorStats { public long numberOfEntriesSinceFirstNotAckedMessage; public int totalNonContiguousDeletedMessagesRange; - public Map consumerRange; - public Map properties; } } From 9050163b68a345c490479d8034106e0bbec2290f Mon Sep 17 00:00:00 2001 From: Marvin Cai Date: Mon, 16 Nov 2020 07:36:35 -0800 Subject: [PATCH 4/5] Update kay hash range string format. --- .../ConsistentHashingStickyKeyConsumerSelector.java | 4 ++-- .../HashRangeAutoSplitStickyKeyConsumerSelector.java | 4 ++-- .../HashRangeExclusiveStickyKeyConsumerSelector.java | 4 ++-- .../broker/service/StickyKeyConsumerSelector.java | 6 +++--- ...rsistentStickyKeyDispatcherMultipleConsumers.java | 4 ++-- .../service/persistent/PersistentSubscription.java | 8 ++++---- ...nsistentHashingStickyKeyConsumerSelectorTest.java | 10 +++++----- ...hRangeAutoSplitStickyKeyConsumerSelectorTest.java | 12 ++++++------ ...hRangeExclusiveStickyKeyConsumerSelectorTest.java | 12 ++++++------ 9 files changed, 32 insertions(+), 32 deletions(-) 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 37049aab1725e..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 @@ -129,7 +129,7 @@ public Consumer select(byte[] stickyKey) { } @Override - public Map> getConsumerRange() { + public Map> getConsumerKeyHashRanges() { Map> result = new LinkedHashMap<>(); rwLock.readLock().lock(); try { @@ -137,7 +137,7 @@ public Map> getConsumerRange() { for (Map.Entry> entry: hashRing.entrySet()) { for (Consumer consumer: entry.getValue()) { result.computeIfAbsent(consumer.consumerName(), key -> new ArrayList<>()) - .add(start + "--" + entry.getKey()); + .add("[" + start + ", " + entry.getKey() + "]"); } start = entry.getKey() + 1; } 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 23de1ff4dd2b2..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 @@ -115,12 +115,12 @@ public Consumer select(byte[] stickyKey) { } @Override - public Map> getConsumerRange() { + 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()); + .add("[" + start + ", " + entry.getKey() + "]"); start = entry.getKey() + 1; } return result; 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 91c260c365d6d..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 @@ -71,7 +71,7 @@ public Consumer select(byte[] stickyKey) { } @Override - public Map> getConsumerRange() { + public Map> getConsumerKeyHashRanges() { Map> result = new HashMap<>(); Map.Entry prev = null; for (Map.Entry entry: rangeMap.entrySet()) { @@ -80,7 +80,7 @@ public Map> getConsumerRange() { } else { if (prev.getValue().equals(entry.getValue())) { result.computeIfAbsent(entry.getValue().consumerName(), key -> new ArrayList<>()) - .add(prev.getKey() + "--" + entry.getKey()); + .add("[" + prev.getKey() + ", " + entry.getKey() + "]"); } prev = null; } 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 3c7f73ebd23b1..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 @@ -48,8 +48,8 @@ public interface StickyKeyConsumerSelector { Consumer select(byte[] stickyKey); /** - * Get range handled by each consumer - * @return A map where key is a consumer name and value is list of range it receiving message for. + * 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> getConsumerRange(); + 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 8dac8889955bc..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,8 +369,8 @@ public LinkedHashMap getRecentlyJoinedConsumers() { return recentlyJoinedConsumers; } - public Map> getConsumerRange() { - return selector.getConsumerRange(); + 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 7d8e40bec4b50..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 @@ -901,8 +901,8 @@ public SubscriptionStats getStats(Boolean getPreciseBacklog) { subStats.lastConsumedFlowTimestamp = lastConsumedFlowTimestamp; Dispatcher dispatcher = this.dispatcher; if (dispatcher != null) { - Map> consumerRanges = getType() == SubType.Key_Shared? - ((PersistentStickyKeyDispatcherMultipleConsumers)dispatcher).getConsumerRange(): null; + Map> consumerKeyHashRanges = getType() == SubType.Key_Shared? + ((PersistentStickyKeyDispatcherMultipleConsumers)dispatcher).getConsumerKeyHashRanges(): null; dispatcher.getConsumers().forEach(consumer -> { ConsumerStats consumerStats = consumer.getStats(); subStats.consumers.add(consumerStats); @@ -915,8 +915,8 @@ 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 (consumerRanges != null && consumerRanges.containsKey(consumer.consumerName())) { - consumerStats.keyHashRange = consumerRanges.get(consumer.consumerName()); + 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 7741cf7dfbe0d..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 @@ -141,7 +141,7 @@ public void testConsumerSelect() throws ConsumerAssignException { @Test - public void testGetConsumerRange() throws BrokerServiceException.ConsumerAssignException { + 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}); @@ -152,10 +152,10 @@ public void testGetConsumerRange() throws BrokerServiceException.ConsumerAssignE } 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.getConsumerRange().entrySet()) { + 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()); } 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 index 6bdab4fb9fd8f..450b92b9fa307 100644 --- 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 @@ -34,7 +34,7 @@ public class HashRangeAutoSplitStickyKeyConsumerSelectorTest { @Test - public void testGetConsumerRange() throws BrokerServiceException.ConsumerAssignException { + 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++) { @@ -44,11 +44,11 @@ public void testGetConsumerRange() throws BrokerServiceException.ConsumerAssignE } 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.getConsumerRange().entrySet()) { + 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()); } 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 f602d294ce9ac..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 @@ -112,7 +112,7 @@ public void testInvalidRangeTotal() { } @Test - public void testGetConsumerRange() throws BrokerServiceException.ConsumerAssignException { + 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}); @@ -130,11 +130,11 @@ public void testGetConsumerRange() throws BrokerServiceException.ConsumerAssignE } 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.getConsumerRange().entrySet()) { + 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()); } From ae082226a2d5995dcdcf4181321ab85f83b8b00c Mon Sep 17 00:00:00 2001 From: Marvin Cai Date: Mon, 16 Nov 2020 08:59:18 -0800 Subject: [PATCH 5/5] Include missing change for ConsumerStats. --- .../org/apache/pulsar/common/policies/data/ConsumerStats.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 a7673d0eadc79..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 @@ -78,8 +78,8 @@ public class ConsumerStats { public long lastAckedTimestamp; public long lastConsumedTimestamp; - /** Hash range assigned to this consumer if is Key_Shared sub mode. **/ - public List keyHashRange; + /** 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;