From 2a05e0a14114dd3a16152b19fc0c6de9fd5f756b Mon Sep 17 00:00:00 2001 From: Masahiro Sakamoto Date: Fri, 3 Jul 2020 21:23:47 +0900 Subject: [PATCH] Consumer is registered on dispatcher even if hash range conflicts on Key_Shared subscription --- .../AbstractDispatcherSingleActiveConsumer.java | 4 ++-- ...PersistentStickyKeyDispatcherMultipleConsumers.java | 10 ++++++++-- ...PersistentStickyKeyDispatcherMultipleConsumers.java | 10 ++++++++-- .../pulsar/client/api/KeySharedSubscriptionTest.java | 10 +++++++++- 4 files changed, 27 insertions(+), 7 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherSingleActiveConsumer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherSingleActiveConsumer.java index 6c5f8a723638b..9948dcc708594 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherSingleActiveConsumer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/AbstractDispatcherSingleActiveConsumer.java @@ -155,8 +155,6 @@ public synchronized void addConsumer(Consumer consumer) throws BrokerServiceExce throw new ConsumerBusyException("Subscription reached max consumers limit"); } - consumers.add(consumer); - if (subscriptionType == SubType.Exclusive && consumer.getKeySharedMeta() != null && consumer.getKeySharedMeta().getHashRangesList() != null @@ -168,6 +166,8 @@ public synchronized void addConsumer(Consumer consumer) throws BrokerServiceExce isKeyHashRangeFiltered = false; } + consumers.add(consumer); + if (!pickAndScheduleActiveConsumer()) { // the active consumer is not changed Consumer currentActiveConsumer = ACTIVE_CONSUMER_UPDATER.get(this); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentStickyKeyDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentStickyKeyDispatcherMultipleConsumers.java index 32cce870608c6..37b29daaf58d9 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentStickyKeyDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/nonpersistent/NonPersistentStickyKeyDispatcherMultipleConsumers.java @@ -47,7 +47,13 @@ public NonPersistentStickyKeyDispatcherMultipleConsumers(NonPersistentTopic topi @Override public synchronized void addConsumer(Consumer consumer) throws BrokerServiceException { super.addConsumer(consumer); - selector.addConsumer(consumer); + try { + selector.addConsumer(consumer); + } catch (BrokerServiceException e) { + consumerSet.removeAll(consumer); + consumerList.remove(consumer); + throw e; + } } @Override @@ -99,4 +105,4 @@ public void sendMessages(List entries) { TOTAL_AVAILABLE_PERMITS_UPDATER.addAndGet(this, -sendMessageInfo.getTotalMessages()); } } -} \ No newline at end of file +} 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 a4a532ff9a3c4..420552c340e37 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 @@ -93,7 +93,13 @@ public class PersistentStickyKeyDispatcherMultipleConsumers extends PersistentDi @Override public synchronized void addConsumer(Consumer consumer) throws BrokerServiceException { super.addConsumer(consumer); - selector.addConsumer(consumer); + try { + selector.addConsumer(consumer); + } catch (BrokerServiceException e) { + consumerSet.removeAll(consumer); + consumerList.remove(consumer); + throw e; + } // If this was the 1st consumer, or if all the messages are already acked, then we // don't need to do anything special @@ -294,4 +300,4 @@ protected Set asyncReplayEntries(Set pos private static final Logger log = LoggerFactory.getLogger(PersistentStickyKeyDispatcherMultipleConsumers.class); -} \ No newline at end of file +} diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java index 511f8c22aa03c..966dfda923cb1 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/KeySharedSubscriptionTest.java @@ -658,7 +658,7 @@ public void testRemoveFirstConsumer() throws Exception { @Test public void testHashRangeConflict() throws PulsarClientException { this.conf.setSubscriptionKeySharedEnable(true); - final String topic = "testHashRangeConflict-" + UUID.randomUUID().toString(); + final String topic = "persistent://public/default/testHashRangeConflict-" + UUID.randomUUID().toString(); final String sub = "test"; Consumer consumer1 = createFixedHashRangesConsumer(topic, sub, Range.of(0,99), Range.of(400, 65535)); @@ -667,6 +667,10 @@ public void testHashRangeConflict() throws PulsarClientException { Consumer consumer2 = createFixedHashRangesConsumer(topic, sub, Range.of(100,399)); Assert.assertTrue(consumer2.isConnected()); + PersistentStickyKeyDispatcherMultipleConsumers dispatcher = (PersistentStickyKeyDispatcherMultipleConsumers) pulsar + .getBrokerService().getTopicReference(topic).get().getSubscription(sub).getDispatcher(); + Assert.assertEquals(dispatcher.getConsumers().size(), 2); + try { createFixedHashRangesConsumer(topic, sub, Range.of(0, 65535)); Assert.fail("Should failed with conflict range."); @@ -679,7 +683,9 @@ public void testHashRangeConflict() throws PulsarClientException { } catch (PulsarClientException.ConsumerAssignException ignore) { } + Assert.assertEquals(dispatcher.getConsumers().size(), 2); consumer1.close(); + Assert.assertEquals(dispatcher.getConsumers().size(), 1); try { createFixedHashRangesConsumer(topic, sub, Range.of(0, 65535)); @@ -705,9 +711,11 @@ public void testHashRangeConflict() throws PulsarClientException { Consumer consumer4 = createFixedHashRangesConsumer(topic, sub, Range.of(50,99)); Assert.assertTrue(consumer4.isConnected()); + Assert.assertEquals(dispatcher.getConsumers().size(), 3); consumer2.close(); consumer3.close(); consumer4.close(); + Assert.assertFalse(dispatcher.isConsumerConnected()); } private Consumer createFixedHashRangesConsumer(String topic, String subscription, Range... ranges) throws PulsarClientException {