diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java index 99971ca809b6b..643e2ad8fe1e9 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentDispatcherMultipleConsumers.java @@ -29,6 +29,7 @@ import java.util.concurrent.CompletableFuture; import java.util.concurrent.CopyOnWriteArrayList; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicIntegerFieldUpdater; import org.apache.bookkeeper.mledger.AsyncCallbacks.ReadEntriesCallback; import org.apache.bookkeeper.mledger.Entry; @@ -81,7 +82,7 @@ public class PersistentDispatcherMultipleConsumers extends AbstractDispatcherMul private Optional delayedDeliveryTracker = Optional.empty(); - protected volatile boolean havePendingRead = false; + protected volatile AtomicBoolean havePendingRead = new AtomicBoolean(false); protected volatile boolean havePendingReplayRead = false; protected boolean shouldRewindBeforeReadingOrReplaying = false; protected final String name; @@ -134,7 +135,7 @@ public synchronized void addConsumer(Consumer consumer) throws BrokerServiceExce return; } if (consumerList.isEmpty()) { - if (havePendingRead || havePendingReplayRead) { + if (havePendingRead.get() || havePendingReplayRead) { // There is a pending read from previous run. We must wait for it to complete and then rewind shouldRewindBeforeReadingOrReplaying = true; } else { @@ -251,12 +252,11 @@ public void readMoreEntries() { } else if (BLOCKED_DISPATCHER_ON_UNACKMSG_UPDATER.get(this) == TRUE) { log.warn("[{}] Dispatcher read is blocked due to unackMessages {} reached to max {}", name, totalUnackedMessages, topic.getMaxUnackedMessagesOnSubscription()); - } else if (!havePendingRead) { + } else if (havePendingRead.compareAndSet(false, true)) { if (log.isDebugEnabled()) { log.debug("[{}] Schedule read of {} messages for {} consumers", name, messagesToRead, consumerList.size()); } - havePendingRead = true; cursor.asyncReadEntriesOrWait(messagesToRead, serviceConfig.getDispatcherMaxReadSizeBytes(), this, ReadType.Normal, topic.getMaxReadPosition()); @@ -402,8 +402,8 @@ public synchronized CompletableFuture disconnectAllConsumers(boolean isRes @Override protected void cancelPendingRead() { - if (havePendingRead && cursor.cancelPendingReadRequest()) { - havePendingRead = false; + if (havePendingRead.get() && cursor.cancelPendingReadRequest()) { + havePendingRead.set(false); } } @@ -432,7 +432,7 @@ public SubType getType() { public synchronized void readEntriesComplete(List entries, Object ctx) { ReadType readType = (ReadType) ctx; if (readType == ReadType.Normal) { - havePendingRead = false; + havePendingRead.set(false); } else { havePendingReplayRead = false; } @@ -594,7 +594,7 @@ public synchronized void readEntriesFailed(ManagedLedgerException exception, Obj } if (readType == ReadType.Normal) { - havePendingRead = false; + havePendingRead.set(false); } else { havePendingReplayRead = false; if (exception instanceof ManagedLedgerException.InvalidReplayPositionException) { @@ -611,7 +611,7 @@ public synchronized void readEntriesFailed(ManagedLedgerException exception, Obj topic.getBrokerService().executor().schedule(() -> { synchronized (PersistentDispatcherMultipleConsumers.this) { - if (!havePendingRead) { + if (!havePendingRead.get()) { log.info("[{}] Retrying read operation", name); readMoreEntries(); } else { @@ -839,7 +839,7 @@ public boolean checkAndUnblockIfStuck() { return false; } // consider dispatch is stuck if : dispatcher has backlog, available-permits and there is no pending read - if (totalAvailablePermits > 0 && !havePendingReplayRead && !havePendingRead + if (totalAvailablePermits > 0 && !havePendingReplayRead && !havePendingRead.get() && cursor.getNumberOfEntriesInBacklog(false) > 0) { log.warn("{}-{} Dispatcher is stuck and unblocking by issuing reads", topic.getName(), name); readMoreEntries(); 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 9340e17aab22c..b9e5b3051c275 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 @@ -61,7 +61,7 @@ public synchronized void readEntryComplete(Entry entry, PendingReadEntryRequest if (ctx.isLast()) { readFailureBackoff.reduceToHalf(); if (readType == ReadType.Normal) { - havePendingRead = false; + havePendingRead.set(false); } else { havePendingReplayRead = false; } @@ -99,11 +99,11 @@ public synchronized void readEntryComplete(Entry entry, PendingReadEntryRequest */ @Override public void canReadMoreEntries(boolean withBackoff) { - havePendingRead = false; + havePendingRead.set(false); topic.getBrokerService().executor().schedule(() -> { topic.getBrokerService().getTopicOrderedExecutor().executeOrdered(topic.getName(), SafeRun.safeRun(() -> { synchronized (PersistentStreamingDispatcherMultipleConsumers.this) { - if (!havePendingRead) { + if (!havePendingRead.get()) { log.info("[{}] Scheduling read operation", name); readMoreEntries(); } else { @@ -129,8 +129,8 @@ public void notifyConsumersEndOfTopic() { @Override protected void cancelPendingRead() { - if (havePendingRead && streamingEntryReader.cancelReadRequests()) { - havePendingRead = false; + if (havePendingRead.get() && streamingEntryReader.cancelReadRequests()) { + havePendingRead.set(false); } } @@ -170,12 +170,11 @@ public void readMoreEntries() { } else if (BLOCKED_DISPATCHER_ON_UNACKMSG_UPDATER.get(this) == TRUE) { log.warn("[{}] Dispatcher read is blocked due to unackMessages {} reached to max {}", name, totalUnackedMessages, topic.getMaxUnackedMessagesOnSubscription()); - } else if (!havePendingRead) { + } else if (havePendingRead.compareAndSet(false, true)) { if (log.isDebugEnabled()) { log.debug("[{}] Schedule read of {} messages for {} consumers", name, messagesToRead, consumerList.size()); } - havePendingRead = true; streamingEntryReader.asyncReadEntries(messagesToRead, serviceConfig.getDispatcherMaxReadSizeBytes(), ReadType.Normal); } else { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java index 972e1a8de49cc..0609e4e290769 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java @@ -1861,7 +1861,7 @@ public CompletableFuture getInternalStats(boolean if (sub.getDispatcher() instanceof PersistentDispatcherMultipleConsumers) { PersistentDispatcherMultipleConsumers dispatcher = (PersistentDispatcherMultipleConsumers) sub .getDispatcher(); - cs.subscriptionHavePendingRead = dispatcher.havePendingRead; + cs.subscriptionHavePendingRead = dispatcher.havePendingRead.get(); cs.subscriptionHavePendingReplayRead = dispatcher.havePendingReplayRead; } else if (sub.getDispatcher() instanceof PersistentDispatcherSingleActiveConsumer) { PersistentDispatcherSingleActiveConsumer dispatcher = (PersistentDispatcherSingleActiveConsumer) sub diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java index a89858ecc7a00..e027f89266442 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/persistent/PersistentTopicTest.java @@ -130,7 +130,7 @@ public void testUnblockStuckSubscription() throws Exception { consumer2.close(); // block sub to read messages - sharedDispatcher.havePendingRead = true; + sharedDispatcher.havePendingRead.set(true); failOverDispatcher.havePendingRead = true; producer.newMessage().value("test").eventTime(5).send(); @@ -146,7 +146,7 @@ public void testUnblockStuckSubscription() throws Exception { assertNull(msg); // allow reads but dispatchers are still blocked - sharedDispatcher.havePendingRead = false; + sharedDispatcher.havePendingRead.set(false); failOverDispatcher.havePendingRead = false; // run task to unblock stuck dispatcher: first iteration sets the lastReadPosition and next iteration will