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..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 @@ -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.getCause(); + } 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(); }