From 46c867dc831ba959e3f235f20591ae38c0b24766 Mon Sep 17 00:00:00 2001 From: nicklixinyang Date: Fri, 6 May 2022 15:24:02 +0800 Subject: [PATCH 1/2] Fix the expired ledger cannot be cleanup when ledger consumed completely --- .../java/org/apache/bookkeeper/mledger/ManagedLedger.java | 5 +++++ .../apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java | 7 +++++-- .../org/apache/pulsar/broker/service/BrokerService.java | 4 +++- .../mledger/offload/jcloud/impl/MockManagedLedger.java | 5 +++++ 4 files changed, 18 insertions(+), 3 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java index 0ebbd514a52bb..c963fd4733ecf 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedLedger.java @@ -620,6 +620,11 @@ void asyncOpenCursor(String name, InitialPosition initialPosition, Map properties, AsyncCallbacks.UpdatePropertiesCallback callback, Object ctx); + /** + * Reset cursor if ledger consumed completely, before trim consumed ledgers in background. + */ + boolean maybeUpdateCursorBeforeTrimmingConsumedLedger(); + /** * Trim consumed ledgers in background. * @param promise diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java index aa7c19b32bd02..38e33ce79e6bb 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java @@ -2296,7 +2296,8 @@ public void addWaitingEntryCallBack(WaitingEntryCallBack cb) { this.waitingEntryCallBacks.add(cb); } - public void maybeUpdateCursorBeforeTrimmingConsumedLedger() { + public boolean maybeUpdateCursorBeforeTrimmingConsumedLedger() { + boolean maybeUpdatedCursor = false; for (ManagedCursor cursor : cursors) { PositionImpl lastAckedPosition = (PositionImpl) cursor.getMarkDeletedPosition(); LedgerInfo currPointedLedger = ledgers.get(lastAckedPosition.getLedgerId()); @@ -2313,13 +2314,14 @@ public void maybeUpdateCursorBeforeTrimmingConsumedLedger() { log.debug("No need to reset cursor: {}, current ledger is the last ledger.", cursor); } } else { - log.warn("Cursor: {} does not exist in the managed-ledger.", cursor); + log.debug("No need to reset cursor: {}, current ledger maybe has removed from managed-ledger.", cursor); } if (!lastAckedPosition.equals((PositionImpl) cursor.getMarkDeletedPosition())) { try { log.info("Reset cursor:{} to {} since ledger consumed completely", cursor, lastAckedPosition); updateCursor((ManagedCursorImpl) cursor, lastAckedPosition); + maybeUpdatedCursor = true; } catch (Exception e) { log.warn("Failed to reset cursor: {} from {} to {}. Trimming thread will retry next time.", cursor, cursor.getMarkDeletedPosition(), lastAckedPosition); @@ -2327,6 +2329,7 @@ public void maybeUpdateCursorBeforeTrimmingConsumedLedger() { } } } + return maybeUpdatedCursor; } private void trimConsumedLedgersInBackground() { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index f19be58f958e1..60fedfb58430e 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -1782,7 +1782,9 @@ private void checkConsumedLedgers() { if (t instanceof PersistentTopic) { Optional.ofNullable(((PersistentTopic) t).getManagedLedger()).ifPresent( managedLedger -> { - managedLedger.trimConsumedLedgersInBackground(Futures.NULL_PROMISE); + if (!managedLedger.maybeUpdateCursorBeforeTrimmingConsumedLedger()){ + managedLedger.trimConsumedLedgersInBackground(Futures.NULL_PROMISE); + } } ); } diff --git a/tiered-storage/jcloud/src/test/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/MockManagedLedger.java b/tiered-storage/jcloud/src/test/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/MockManagedLedger.java index 8ababead6306c..3b44fdfe46dbd 100644 --- a/tiered-storage/jcloud/src/test/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/MockManagedLedger.java +++ b/tiered-storage/jcloud/src/test/java/org/apache/bookkeeper/mledger/offload/jcloud/impl/MockManagedLedger.java @@ -328,6 +328,11 @@ public void asyncSetProperties(Map properties, } + @Override + public boolean maybeUpdateCursorBeforeTrimmingConsumedLedger() { + return false; + } + @Override public void trimConsumedLedgersInBackground(CompletableFuture promise) { From 326dad17f77841a8b8f9cc2d08c3e80b90285a8b Mon Sep 17 00:00:00 2001 From: nicklixinyang Date: Fri, 5 Aug 2022 14:17:46 +0800 Subject: [PATCH 2/2] add unit test and comment for clean expired ledgers after expired messages checked --- .../pulsar/broker/service/BrokerService.java | 2 + .../service/ConsumedLedgersTrimTest.java | 66 +++++++++++++++++++ 2 files changed, 68 insertions(+) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java index 60fedfb58430e..5b3e2083d37dc 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/BrokerService.java @@ -1782,6 +1782,8 @@ private void checkConsumedLedgers() { if (t instanceof PersistentTopic) { Optional.ofNullable(((PersistentTopic) t).getManagedLedger()).ifPresent( managedLedger -> { + // After update cursor, trimConsumedLedgersInBackground will be invoked, + // avoid invoke repeatedly here. if (!managedLedger.maybeUpdateCursorBeforeTrimmingConsumedLedger()){ managedLedger.trimConsumedLedgersInBackground(Futures.NULL_PROMISE); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ConsumedLedgersTrimTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ConsumedLedgersTrimTest.java index 355036bdb25ce..1a6c1576b336f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ConsumedLedgersTrimTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/ConsumedLedgersTrimTest.java @@ -22,16 +22,21 @@ import static org.testng.Assert.assertNotEquals; import static org.testng.Assert.assertNotNull; +import java.lang.reflect.Method; import java.util.concurrent.CompletableFuture; import java.util.concurrent.TimeUnit; import lombok.Cleanup; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; +import org.apache.bookkeeper.mledger.Position; +import org.apache.bookkeeper.mledger.impl.ManagedCursorImpl; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; +import org.apache.pulsar.broker.service.persistent.PersistentMessageExpiryMonitor; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.Message; import org.apache.pulsar.client.api.MessageId; import org.apache.pulsar.client.api.Producer; +import org.awaitility.Awaitility; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import org.testng.Assert; @@ -179,4 +184,65 @@ public void testConsumedLedgersTrimNoSubscriptions() throws Exception { assertEquals(messageIdAfterTrim, MessageId.earliest); } + + @Test + public void testExpiredLedgerDeletionAfterExpiredMessageChecked() throws Exception { + super.baseSetup(); + final String ledgerAndCursorName = "testExpiredLedgerDeletionAfterExpiredMessageChecked"; + final int totalEntries = 10; + final int ttlSeconds = 1; + + final String topicName = "persistent://prop/ns-abc/testExpiredLedgerDeletionAfterExpiredMessageChecked"; + + @Cleanup + Producer producer = pulsarClient.newProducer() + .topic(topicName) + .producerName("producer-name") + .create(); + + //set retention parameters, the ledgers are to be deleted as soon as possible + PersistentTopic persistentTopic = (PersistentTopic) pulsar.getBrokerService().getOrCreateTopic(topicName).get(); + ManagedLedgerConfig managedLedgerConfig = persistentTopic.getManagedLedger().getConfig(); + managedLedgerConfig.setRetentionSizeInMB(10); + managedLedgerConfig.setRetentionTime(1, TimeUnit.SECONDS); + managedLedgerConfig.setMaxEntriesPerLedger(1000); + managedLedgerConfig.setMinimumRolloverTime(1, TimeUnit.MILLISECONDS); + managedLedgerConfig.setMaximumRolloverTime(50, TimeUnit.MILLISECONDS); + + //create managedLedger and managedCursor + ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) persistentTopic.getManagedLedger(); + ManagedCursorImpl managedCursor = (ManagedCursorImpl) managedLedger.openCursor(ledgerAndCursorName); + + // write some messages + for (int i = 0; i < totalEntries; i++) { + producer.send(("msg" + i).getBytes()); + } + + //make sure that all entries should be deleted + Thread.sleep(TimeUnit.SECONDS.toMillis(ttlSeconds)); + managedLedger.rollCurrentLedgerIfFull(); + + PersistentMessageExpiryMonitor monitor = new PersistentMessageExpiryMonitor(topicName, managedCursor.getName(), managedCursor, null); + Position previousMarkDelete = null; + for (int i = 0; i < totalEntries; i++) { + monitor.expireMessages(1); + Position previousPos = previousMarkDelete; + retryStrategically( + (test) -> managedCursor.getMarkDeletedPosition() != null && !managedCursor.getMarkDeletedPosition().equals(previousPos), + 5, 100); + previousMarkDelete = managedCursor.getMarkDeletedPosition(); + } + + //invoke checkConsumedLedgers() to clean the expired ledgers + Method checkConsumedLedgers = BrokerService.class.getDeclaredMethod("checkConsumedLedgers"); + checkConsumedLedgers.setAccessible(true); + checkConsumedLedgers.invoke(pulsar.getBrokerService()); + + Awaitility.await().untilAsserted(() + -> assertEquals(managedLedger.getLedgersInfo().size(), 1)); + assertEquals(managedLedger.getLedgersInfo().firstEntry().getValue().getEntries(), 0); + + managedCursor.close(); + managedLedger.close(); + } }