From cc93e09b980b85eefd64bc0181c26692c9290b4c Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Thu, 20 Oct 2022 00:24:11 +0800 Subject: [PATCH 1/2] [broker] Avoid put all delayed messages to redeliver queue --- .../mledger/ManagedLedgerException.java | 18 ++++++++++++++++++ .../mledger/impl/ManagedCursorImpl.java | 6 ++++-- ...PersistentDispatcherMultipleConsumers.java | 19 ++++++++++++++----- 3 files changed, 36 insertions(+), 7 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerException.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerException.java index fa6dc59d14753..3329e2bc17a40 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerException.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedgerException.java @@ -18,6 +18,8 @@ */ package org.apache.bookkeeper.mledger; +import java.util.Set; +import lombok.Getter; import org.apache.bookkeeper.common.annotation.InterfaceAudience; import org.apache.bookkeeper.common.annotation.InterfaceStability; @@ -192,6 +194,22 @@ public ManagedLedgerFactoryClosedException(Throwable e) { } } + + public static class CursorReplyFailedException extends ManagedLedgerException { + @Getter + private final Set replyPositions; + + public CursorReplyFailedException(ManagedLedgerException ex, Set replyPositions) { + super(ex); + this.replyPositions = replyPositions; + } + + @Override + public synchronized ManagedLedgerException getCause() { + return (ManagedLedgerException) super.getCause(); + } + } + public static class ConcurrentWaitCallbackException extends ManagedLedgerException { public ConcurrentWaitCallbackException() { diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java index 20f760d893adc..d7c3de439fbdf 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedCursorImpl.java @@ -1435,7 +1435,8 @@ public synchronized void readEntryComplete(Entry entry, Object ctx) { // and not add it to the list entry.release(); if (--pendingCallbacks == 0) { - callback.readEntriesFailed(exception.get(), ctx); + callback.readEntriesFailed( + new ManagedLedgerException.CursorReplyFailedException(exception.get(), positions), ctx); } } else { entries.add(entry); @@ -1456,7 +1457,8 @@ public synchronized void readEntryFailed(ManagedLedgerException mle, Object ctx) entries.forEach(Entry::release); } if (--pendingCallbacks == 0) { - callback.readEntriesFailed(exception.get(), ctx); + callback.readEntriesFailed( + new ManagedLedgerException.CursorReplyFailedException(exception.get(), positions), ctx); } } }; 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 9f09e60abb29e..0adacf68c1a80 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 @@ -767,10 +767,22 @@ private boolean sendChunkedMessagesToConsumers(ReadType readType, } @Override - public synchronized void readEntriesFailed(ManagedLedgerException exception, Object ctx) { + public synchronized void readEntriesFailed(ManagedLedgerException rawException, Object ctx) { ReadType readType = (ReadType) ctx; long waitTimeMillis = readFailureBackoff.next(); + final ManagedLedgerException exception; + + // Add failed reply messages back to redelivery messages. + if (rawException instanceof ManagedLedgerException.CursorReplyFailedException) { + ManagedLedgerException.CursorReplyFailedException cursorReplyFailedException = + (ManagedLedgerException.CursorReplyFailedException) rawException; + cursorReplyFailedException.getReplyPositions().forEach(replyPosition -> + redeliveryMessages.add(replyPosition.getLedgerId(), replyPosition.getEntryId())); + exception = cursorReplyFailedException; + } else { + exception = rawException; + } if (exception instanceof NoMoreEntriesToReadException) { if (cursor.getNumberOfEntriesInBacklog(false) == 0) { @@ -1025,10 +1037,7 @@ protected synchronized Set getMessagesToReplayNow(int maxMessagesT return redeliveryMessages.getMessagesToReplayNow(maxMessagesToRead); } else if (delayedDeliveryTracker.isPresent() && delayedDeliveryTracker.get().hasMessageAvailable()) { delayedDeliveryTracker.get().resetTickTime(topic.getDelayedDeliveryTickTimeMillis()); - Set messagesAvailableNow = - delayedDeliveryTracker.get().getScheduledMessages(maxMessagesToRead); - messagesAvailableNow.forEach(p -> redeliveryMessages.add(p.getLedgerId(), p.getEntryId())); - return messagesAvailableNow; + return delayedDeliveryTracker.get().getScheduledMessages(maxMessagesToRead); } else { return Collections.emptySet(); } From 25c8d917ad9402ae97c6d7bf9f75be04c895538a Mon Sep 17 00:00:00 2001 From: mattisonchao Date: Thu, 20 Oct 2022 00:33:19 +0800 Subject: [PATCH 2/2] Fix exception --- .../persistent/PersistentDispatcherMultipleConsumers.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 0adacf68c1a80..e2759a795f4e8 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 @@ -779,7 +779,7 @@ public synchronized void readEntriesFailed(ManagedLedgerException rawException, (ManagedLedgerException.CursorReplyFailedException) rawException; cursorReplyFailedException.getReplyPositions().forEach(replyPosition -> redeliveryMessages.add(replyPosition.getLedgerId(), replyPosition.getEntryId())); - exception = cursorReplyFailedException; + exception = cursorReplyFailedException.getCause(); } else { exception = rawException; }