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(); } }