From 9b22af984a914e61daaf462592564adada78119b Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Wed, 2 Nov 2022 00:42:02 +0800 Subject: [PATCH 1/2] [fix][test]fix flaky test SimpleProducerConsumerTestStreamingDispatcherTest.testSharedSamePriorityConsumer --- ...entStreamingDispatcherMultipleConsumers.java | 17 +++++++++++++++-- .../client/api/SimpleProducerConsumerTest.java | 6 ++++++ 2 files changed, 21 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStreamingDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStreamingDispatcherMultipleConsumers.java index 5643f91d93d63..66f2473a639e2 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStreamingDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentStreamingDispatcherMultipleConsumers.java @@ -20,6 +20,7 @@ import static org.apache.bookkeeper.mledger.util.SafeRun.safeRun; import com.google.common.collect.Lists; +import java.util.List; import java.util.Set; import java.util.concurrent.TimeUnit; import lombok.extern.slf4j.Slf4j; @@ -52,6 +53,18 @@ public PersistentStreamingDispatcherMultipleConsumers(PersistentTopic topic, Man super(topic, cursor, subscription); } + protected final synchronized boolean sendMessagesToConsumers(ReadType readType, List entries, + boolean isLastEntryInBatch) { + sendInProgress = true; + try { + return trySendMessagesToConsumers(readType, entries); + } finally { + if (isLastEntryInBatch) { + sendInProgress = false; + } + } + } + /** * {@inheritDoc} */ @@ -100,12 +113,12 @@ public synchronized void readEntryComplete(Entry entry, PendingReadEntryRequest // in a separate thread, and we want to prevent more reads sendInProgress = true; dispatchMessagesThread.execute(safeRun(() -> { - if (sendMessagesToConsumers(readType, Lists.newArrayList(entry))) { + if (sendMessagesToConsumers(readType, Lists.newArrayList(entry), ctx.isLast())) { readMoreEntries(); } })); } else { - if (sendMessagesToConsumers(readType, Lists.newArrayList(entry))) { + if (sendMessagesToConsumers(readType, Lists.newArrayList(entry), ctx.isLast())) { readMoreEntriesAsync(); } } 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 4329083b30c9d..7f868be4a72e1 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 @@ -2411,6 +2411,7 @@ public void testSharedSamePriorityConsumer() throws Exception { PulsarClient newPulsarClient = newPulsarClient(lookupUrl.toString(), 0);// Creates new client connection Consumer c1 = newPulsarClient.newConsumer() .topic("persistent://my-property/my-ns/my-topic2").subscriptionName("my-subscriber-name") + .consumerName("c1") .subscriptionType(SubscriptionType.Shared).receiverQueueSize(queueSize) .acknowledgmentGroupTime(0, TimeUnit.SECONDS).subscribe(); @@ -2418,6 +2419,7 @@ public void testSharedSamePriorityConsumer() throws Exception { PulsarClient newPulsarClient1 = newPulsarClient(lookupUrl.toString(), 0);// Creates new client connection Consumer c2 = newPulsarClient1.newConsumer() .topic("persistent://my-property/my-ns/my-topic2").subscriptionName("my-subscriber-name") + .consumerName("c1") .subscriptionType(SubscriptionType.Shared).receiverQueueSize(queueSize) .acknowledgmentGroupTime(0, TimeUnit.SECONDS).subscribe(); List> futures = new ArrayList<>(); @@ -2462,6 +2464,7 @@ public void testSharedSamePriorityConsumer() throws Exception { PulsarClient newPulsarClient2 = newPulsarClient(lookupUrl.toString(), 0);// Creates new client connection Consumer c3 = newPulsarClient2.newConsumer() .topic("persistent://my-property/my-ns/my-topic2").subscriptionName("my-subscriber-name") + .consumerName("c3") .subscriptionType(SubscriptionType.Shared).receiverQueueSize(queueSize) .acknowledgmentGroupTime(0, TimeUnit.SECONDS).subscribe(); @@ -2469,6 +2472,7 @@ public void testSharedSamePriorityConsumer() throws Exception { PulsarClient newPulsarClient3 = newPulsarClient(lookupUrl.toString(), 0);// Creates new client connection Consumer c4 = newPulsarClient3.newConsumer() .topic("persistent://my-property/my-ns/my-topic2").subscriptionName("my-subscriber-name") + .consumerName("c4") .subscriptionType(SubscriptionType.Shared).receiverQueueSize(queueSize) .acknowledgmentGroupTime(0, TimeUnit.SECONDS).subscribe(); @@ -2476,6 +2480,7 @@ public void testSharedSamePriorityConsumer() throws Exception { PulsarClient newPulsarClient4 = newPulsarClient(lookupUrl.toString(), 0);// Creates new client connection Consumer c5 = newPulsarClient4.newConsumer() .topic("persistent://my-property/my-ns/my-topic2").subscriptionName("my-subscriber-name") + .consumerName("c5") .subscriptionType(SubscriptionType.Shared).receiverQueueSize(queueSize) .acknowledgmentGroupTime(0, TimeUnit.SECONDS).subscribe(); @@ -2520,6 +2525,7 @@ public void testSharedSamePriorityConsumer() throws Exception { c4.close(); c5.close(); pulsar.getConfiguration().setMaxUnackedMessagesPerConsumer(maxUnAckMsgs); + admin.topics().delete("persistent://my-property/my-ns/my-topic2", false); log.info("-- Exiting {} test --", methodName); } From 11f101a127c0d03fe65952585908cb75e96cda99 Mon Sep 17 00:00:00 2001 From: fengyubiao Date: Wed, 2 Nov 2022 11:23:21 +0800 Subject: [PATCH 2/2] remove unnecessary change --- .../pulsar/client/api/SimpleProducerConsumerTest.java | 6 ------ 1 file changed, 6 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 7f868be4a72e1..4329083b30c9d 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 @@ -2411,7 +2411,6 @@ public void testSharedSamePriorityConsumer() throws Exception { PulsarClient newPulsarClient = newPulsarClient(lookupUrl.toString(), 0);// Creates new client connection Consumer c1 = newPulsarClient.newConsumer() .topic("persistent://my-property/my-ns/my-topic2").subscriptionName("my-subscriber-name") - .consumerName("c1") .subscriptionType(SubscriptionType.Shared).receiverQueueSize(queueSize) .acknowledgmentGroupTime(0, TimeUnit.SECONDS).subscribe(); @@ -2419,7 +2418,6 @@ public void testSharedSamePriorityConsumer() throws Exception { PulsarClient newPulsarClient1 = newPulsarClient(lookupUrl.toString(), 0);// Creates new client connection Consumer c2 = newPulsarClient1.newConsumer() .topic("persistent://my-property/my-ns/my-topic2").subscriptionName("my-subscriber-name") - .consumerName("c1") .subscriptionType(SubscriptionType.Shared).receiverQueueSize(queueSize) .acknowledgmentGroupTime(0, TimeUnit.SECONDS).subscribe(); List> futures = new ArrayList<>(); @@ -2464,7 +2462,6 @@ public void testSharedSamePriorityConsumer() throws Exception { PulsarClient newPulsarClient2 = newPulsarClient(lookupUrl.toString(), 0);// Creates new client connection Consumer c3 = newPulsarClient2.newConsumer() .topic("persistent://my-property/my-ns/my-topic2").subscriptionName("my-subscriber-name") - .consumerName("c3") .subscriptionType(SubscriptionType.Shared).receiverQueueSize(queueSize) .acknowledgmentGroupTime(0, TimeUnit.SECONDS).subscribe(); @@ -2472,7 +2469,6 @@ public void testSharedSamePriorityConsumer() throws Exception { PulsarClient newPulsarClient3 = newPulsarClient(lookupUrl.toString(), 0);// Creates new client connection Consumer c4 = newPulsarClient3.newConsumer() .topic("persistent://my-property/my-ns/my-topic2").subscriptionName("my-subscriber-name") - .consumerName("c4") .subscriptionType(SubscriptionType.Shared).receiverQueueSize(queueSize) .acknowledgmentGroupTime(0, TimeUnit.SECONDS).subscribe(); @@ -2480,7 +2476,6 @@ public void testSharedSamePriorityConsumer() throws Exception { PulsarClient newPulsarClient4 = newPulsarClient(lookupUrl.toString(), 0);// Creates new client connection Consumer c5 = newPulsarClient4.newConsumer() .topic("persistent://my-property/my-ns/my-topic2").subscriptionName("my-subscriber-name") - .consumerName("c5") .subscriptionType(SubscriptionType.Shared).receiverQueueSize(queueSize) .acknowledgmentGroupTime(0, TimeUnit.SECONDS).subscribe(); @@ -2525,7 +2520,6 @@ public void testSharedSamePriorityConsumer() throws Exception { c4.close(); c5.close(); pulsar.getConfiguration().setMaxUnackedMessagesPerConsumer(maxUnAckMsgs); - admin.topics().delete("persistent://my-property/my-ns/my-topic2", false); log.info("-- Exiting {} test --", methodName); }