From bd6966ec1fc83f9ed5a14f3c9362ec1e6f4741d9 Mon Sep 17 00:00:00 2001 From: aloyszhang Date: Sat, 7 Jan 2023 10:39:30 +0800 Subject: [PATCH] fix ttl expiration block due to no-recoverable exception when autoSkipNonRecoverableData=true --- .../broker/service/persistent/PersistentTopic.java | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) 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();