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 baca6bf078cf0..ec4d400d715eb 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 @@ -64,6 +64,7 @@ public abstract class AbstractDispatcherSingleActiveConsumer extends AbstractBas protected boolean isFirstRead = true; private static final int CONSUMER_CONSISTENT_HASH_REPLICAS = 100; + private int addConsumerFailedAttemptCount = 0; public AbstractDispatcherSingleActiveConsumer(SubType subscriptionType, int partitionIndex, String topicName, Subscription subscription, @@ -176,7 +177,7 @@ public synchronized CompletableFuture addConsumer(Consumer consumer) { return FutureUtil.failedFuture(new ConsumerBusyException("Exclusive consumer is already" + " connected")); } else { - return addConsumer(consumer); + return internalAddConsumer(actConsumer, consumer); } }); } else { @@ -225,6 +226,25 @@ public synchronized CompletableFuture addConsumer(Consumer consumer) { return CompletableFuture.completedFuture(null); } + /** + * This method is used to help debugging addConsumer failed for exclusive subscription. + * @param actConsumer + * @param consumer + * @return + */ + private synchronized CompletableFuture internalAddConsumer(Consumer actConsumer, Consumer consumer) { + addConsumerFailedAttemptCount++; + if (addConsumerFailedAttemptCount >= 5) { + log.warn("Added consumer failed with attempt count {}, consumers: {}, active consumer: {}," + + "active state : {}", + addConsumerFailedAttemptCount, consumers, actConsumer, actConsumer.cnx().isActive()); + addConsumerFailedAttemptCount = 0; + return FutureUtil.failedFuture(new ConsumerBusyException("Exclusive consumer is already" + + " connected")); + } + return addConsumer(consumer); + } + public synchronized void removeConsumer(Consumer consumer) throws BrokerServiceException { log.info("Removing consumer {}", consumer); if (!consumers.remove(consumer)) {