diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java index c45c2b1522f59..8830d14656dc2 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/client/api/SimpleProducerConsumerTest.java @@ -84,6 +84,7 @@ import org.apache.commons.lang3.reflect.FieldUtils; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.PulsarVersion; +import org.apache.pulsar.broker.BrokerTestUtil; import org.apache.pulsar.broker.PulsarService; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.client.admin.PulsarAdminException; @@ -107,7 +108,9 @@ import org.apache.pulsar.common.compression.CompressionCodec; import org.apache.pulsar.common.compression.CompressionCodecProvider; import org.apache.pulsar.common.naming.TopicName; +import org.apache.pulsar.common.policies.data.ConsumerStats; import org.apache.pulsar.common.policies.data.PublisherStats; +import org.apache.pulsar.common.policies.data.SubscriptionStats; import org.apache.pulsar.common.policies.data.TopicStats; import org.apache.pulsar.common.protocol.Commands; import org.apache.pulsar.common.schema.SchemaType; @@ -190,6 +193,16 @@ public Object[][] ackReceiptEnabledAndSubscriptionTypes() { }; } + @DataProvider(name = "subscriptionTypes") + public Object[][] subType() { + return new Object[][] { + {SubscriptionType.Shared}, + {SubscriptionType.Key_Shared}, + {SubscriptionType.Exclusive}, + {SubscriptionType.Failover} + }; + } + @AfterClass(alwaysRun = true) @Override protected void cleanup() throws Exception { @@ -340,6 +353,58 @@ public Object[][] codecProvider() { return new Object[][] { { 0 }, { 1000 } }; } + @Test(dataProvider = "subscriptionTypes") + public void testConsumerReconnectTwice(SubscriptionType subscriptionType) throws Exception { + final String topicName = BrokerTestUtil.newUniqueName("persistent://my-property/my-ns/tp_"); + final String subscriptionName = "subscription1"; + admin.topics().createNonPartitionedTopic(topicName); + admin.topics().createSubscription(topicName, subscriptionName, MessageId.earliest); + // Create producer and consumer. + ConsumerImpl consumer = (ConsumerImpl) pulsarClient.newConsumer(Schema.STRING) + .subscriptionType(subscriptionType) + .receiverQueueSize(1000).topic(topicName).subscriptionName(subscriptionName).subscribe(); + Producer producer = pulsarClient.newProducer(Schema.STRING).enableBatching(false) + .topic(topicName).create(); + int sendMessageCount = 10; + for (int i = 0; i < sendMessageCount; i++){ + producer.send("msg- " + i); + } + Awaitility.await().untilAsserted(() -> { + assertEquals(consumer.numMessagesInQueue(), sendMessageCount); + }); + printConsumerStats(topicName, subscriptionName); + + // Do the second subscribe. + consumer.connectionOpened(consumer.getClientCnx()); + + // Verify messages are not lost. + List> messages = new ArrayList<>(); + while (true) { + Message message = consumer.receive(2, TimeUnit.SECONDS); + if (message == null) { + break; + } + messages.add(message); + } + printConsumerStats(topicName, subscriptionName); + assertEquals(messages.size(), sendMessageCount); + + // cleanup. + consumer.close(); + producer.close(); + admin.topics().delete(topicName, false); + } + + private void printConsumerStats(String topicName, String subscriptionName) throws Exception { + SubscriptionStats subscriptionStats = + admin.topics().getStats(topicName).getSubscriptions().get(subscriptionName); + ConsumerStats consumerStats = + admin.topics().getStats(topicName).getSubscriptions().get(subscriptionName).getConsumers().get(0); + log.info("msgBacklog: {}, msgOutCounter: {}, unackedMessages: {}, availablePermits: {}", + subscriptionStats.getMsgBacklog(), consumerStats.getMsgOutCounter(), + consumerStats.getUnackedMessages(), consumerStats.getAvailablePermits()); + } + @Test(timeOut = 100000, dataProvider = "batch") public void testSyncProducerAndConsumer(int batchMessageDelayMs) throws Exception { log.info("-- Starting {} test --", methodName); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java index 4a84e765065f2..3b287468acdf7 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ConsumerImpl.java @@ -136,6 +136,7 @@ public class ConsumerImpl extends ConsumerBase implements ConnectionHandle protected volatile MessageId lastDequeuedMessageId = MessageId.earliest; private volatile MessageId lastMessageIdInBroker = MessageId.earliest; + protected volatile CompletableFuture inProgressSubscribeFuture; private final long lookupDeadline; @@ -761,6 +762,18 @@ public void negativeAcknowledge(Message message) { @Override public void connectionOpened(final ClientCnx cnx) { + synchronized (this) { + // Wait the previous subscribe done. + if (inProgressSubscribeFuture != null && !inProgressSubscribeFuture.isDone()){ + return; + } + // If success. + if (getState() == State.Ready) { + return; + } + // Do subscribe if previous was failed. + inProgressSubscribeFuture = new CompletableFuture<>(); + } previousExceptions.clear(); if (getState() == State.Closing || getState() == State.Closed) { @@ -906,6 +919,12 @@ public void connectionOpened(final ClientCnx cnx) { reconnectLater(e.getCause()); } return null; + }).whenComplete((ignore, ex) -> { + if (ex == null) { + inProgressSubscribeFuture.complete(null); + } else { + inProgressSubscribeFuture.completeExceptionally(ex); + } }); } }