From 9bbfacbd08a5d8469a49eaa61e77a29b8f5a4c06 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Fri, 16 Jun 2023 13:58:53 +0800 Subject: [PATCH 1/6] [fix] [client] Messages lost when consumer tries connect twice --- .../api/SimpleProducerConsumerTest.java | 27 +++++++++++++++++++ .../pulsar/client/impl/ConsumerImpl.java | 19 +++++++++++++ 2 files changed, 46 insertions(+) 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..2bb3a4047f314 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; @@ -340,6 +341,32 @@ public Object[][] codecProvider() { return new Object[][] { { 0 }, { 1000 } }; } + @Test + public void testConsumerReconnectTwice() 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); + + ConsumerImpl consumer = (ConsumerImpl) pulsarClient.newConsumer(Schema.STRING) + .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); + }); + + consumer.connectionOpened(consumer.getClientCnx()); + Awaitility.await().untilAsserted(() -> { + assertEquals(consumer.numMessagesInQueue(), sendMessageCount); + }); + } + @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..92c9307566b4f 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.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); + } }); } } From 8248d75e0b99c7fb86583c1107814ca490536ec0 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Fri, 16 Jun 2023 14:00:39 +0800 Subject: [PATCH 2/6] [fix] [client] Messages lost when consumer tries connect twice --- .../main/java/org/apache/pulsar/client/impl/ConsumerImpl.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 92c9307566b4f..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 @@ -764,7 +764,7 @@ public void negativeAcknowledge(Message message) { public void connectionOpened(final ClientCnx cnx) { synchronized (this) { // Wait the previous subscribe done. - if (!inProgressSubscribeFuture.isDone()){ + if (inProgressSubscribeFuture != null && !inProgressSubscribeFuture.isDone()){ return; } // If success. From 2a7a0bdbe87b2af293a4636ff4ff54213b6a27ca Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Fri, 16 Jun 2023 14:22:24 +0800 Subject: [PATCH 3/6] improve test --- .../pulsar/client/api/SimpleProducerConsumerTest.java | 8 ++++++++ 1 file changed, 8 insertions(+) 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 2bb3a4047f314..fce1d843d66e8 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 @@ -108,6 +108,7 @@ 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.TopicStats; import org.apache.pulsar.common.protocol.Commands; @@ -363,8 +364,15 @@ public void testConsumerReconnectTwice() throws Exception { consumer.connectionOpened(consumer.getClientCnx()); Awaitility.await().untilAsserted(() -> { + ConsumerStats consumerStats = + admin.topics().getStats(topicName).getSubscriptions().get(subscriptionName).getConsumers().get(0); + log.info("consumerStats: " + consumerStats.toString()); assertEquals(consumer.numMessagesInQueue(), sendMessageCount); }); + + consumer.close(); + producer.close(); + admin.topics().delete(topicName, false); } @Test(timeOut = 100000, dataProvider = "batch") From 696314ace27d893f8dc2017f9a046db62c6c82ef Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Fri, 16 Jun 2023 14:34:09 +0800 Subject: [PATCH 4/6] improve test --- .../api/SimpleProducerConsumerTest.java | 35 ++++++++++++++----- 1 file changed, 27 insertions(+), 8 deletions(-) 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 fce1d843d66e8..2baad768ca28b 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 @@ -110,6 +110,7 @@ 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; @@ -348,12 +349,11 @@ public void testConsumerReconnectTwice() throws Exception { 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) .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); @@ -361,20 +361,39 @@ public void testConsumerReconnectTwice() throws Exception { Awaitility.await().untilAsserted(() -> { assertEquals(consumer.numMessagesInQueue(), sendMessageCount); }); + printConsumerStats(topicName, subscriptionName); + // Do the second subscribe. consumer.connectionOpened(consumer.getClientCnx()); - Awaitility.await().untilAsserted(() -> { - ConsumerStats consumerStats = - admin.topics().getStats(topicName).getSubscriptions().get(subscriptionName).getConsumers().get(0); - log.info("consumerStats: " + consumerStats.toString()); - assertEquals(consumer.numMessagesInQueue(), sendMessageCount); - }); + // 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); From 212a11614329ba6c468b4e35f4f2ee1e9c800321 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Fri, 16 Jun 2023 17:45:24 +0800 Subject: [PATCH 5/6] improve test --- .../org/apache/pulsar/client/api/SimpleProducerConsumerTest.java | 1 + 1 file changed, 1 insertion(+) 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 2baad768ca28b..80b8d5d6d5521 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 @@ -351,6 +351,7 @@ public void testConsumerReconnectTwice() throws Exception { admin.topics().createSubscription(topicName, subscriptionName, MessageId.earliest); // Create producer and consumer. ConsumerImpl consumer = (ConsumerImpl) pulsarClient.newConsumer(Schema.STRING) + .subscriptionType(SubscriptionType.Shared) .receiverQueueSize(1000).topic(topicName).subscriptionName(subscriptionName).subscribe(); Producer producer = pulsarClient.newProducer(Schema.STRING).enableBatching(false) .topic(topicName).create(); From d583dc2d2030b39347e8d74ad11b276134e0c3d9 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Fri, 16 Jun 2023 17:57:25 +0800 Subject: [PATCH 6/6] improve test --- .../client/api/SimpleProducerConsumerTest.java | 16 +++++++++++++--- 1 file changed, 13 insertions(+), 3 deletions(-) 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 80b8d5d6d5521..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 @@ -193,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 { @@ -343,15 +353,15 @@ public Object[][] codecProvider() { return new Object[][] { { 0 }, { 1000 } }; } - @Test - public void testConsumerReconnectTwice() throws Exception { + @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.Shared) + .subscriptionType(subscriptionType) .receiverQueueSize(1000).topic(topicName).subscriptionName(subscriptionName).subscribe(); Producer producer = pulsarClient.newProducer(Schema.STRING).enableBatching(false) .topic(topicName).create();