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 42f73a7032835..6abaa19f85426 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 @@ -67,6 +67,7 @@ import org.apache.bookkeeper.mledger.ManagedLedgerException.ManagedLedgerFencedException; import org.apache.bookkeeper.mledger.ManagedLedgerException.ManagedLedgerTerminatedException; import org.apache.bookkeeper.mledger.ManagedLedgerException.MetadataNotFoundException; +import org.apache.bookkeeper.mledger.ManagedLedgerException.NonRecoverableLedgerException; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.impl.ManagedCursorContainer; import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl; @@ -2853,7 +2854,14 @@ public boolean isOldestMessageExpired(ManagedCursor cursor, int messageTTLInSeco (int) (messageTTLInSeconds * MESSAGE_EXPIRY_THRESHOLD), entryTimestamp); } } catch (Exception e) { - log.warn("[{}] Error while getting the oldest message", topic, e); + if (brokerService.pulsar().getConfiguration().isAutoSkipNonRecoverableData() + && e instanceof NonRecoverableLedgerException) { + // NonRecoverableLedgerException means the ledger or entry can't be read anymore. + // if AutoSkipNonRecoverableData is set to true, just return true here. + return true; + } else { + log.warn("[{}] Error while getting the oldest message", topic, e); + } } finally { if (entry != null) { entry.release();