From 5993dca35fd9f0094ab0ecdb4eaf51a30ce037b6 Mon Sep 17 00:00:00 2001 From: wuzhanpeng Date: Tue, 30 Mar 2021 18:37:55 +0800 Subject: [PATCH 1/5] bugfix for wrong timeunit in updating lastLedgerCreationInitiationTimestamp --- .../mledger/impl/ManagedLedgerImpl.java | 47 +++++++++++++++---- .../mledger/impl/ManagedLedgerTest.java | 46 +++++++++++++++++- .../CurrentLedgerRolloverIfFullTest.java | 3 ++ 3 files changed, 85 insertions(+), 11 deletions(-) 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 e4c4b4a75c11a..08f4420213244 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 @@ -1388,8 +1388,6 @@ public synchronized void createComplete(int rc, final LedgerHandle lh, Object ct } else { log.info("[{}] Created new ledger {}", name, lh.getId()); ledgers.put(lh.getId(), LedgerInfo.newBuilder().setLedgerId(lh.getId()).setTimestamp(0).build()); - final long previousEntries = currentLedgerEntries; - final long previousLedgerId = currentLedger.getId(); currentLedger = lh; currentLedgerEntries = 0; currentLedgerSize = 0; @@ -1407,14 +1405,6 @@ public void operationComplete(Void v, Stat stat) { mbean.addLedgerSwitchLatencySample(System.currentTimeMillis() - lastLedgerCreationInitiationTimestamp, TimeUnit.MILLISECONDS); } - // Move cursor read point to new ledger - for (ManagedCursor cursor : cursors) { - PositionImpl markDeletedPosition = (PositionImpl) cursor.getMarkDeletedPosition(); - if (markDeletedPosition.getLedgerId() == previousLedgerId && markDeletedPosition.getEntryId() + 1 >= previousEntries) { - // All entries in last ledger are marked delete, move read point to the new ledger - updateCursor((ManagedCursorImpl) cursor, PositionImpl.get(currentLedger.getId(), -1)); - } - } } @Override @@ -2140,6 +2130,40 @@ public void addWaitingEntryCallBack(WaitingEntryCallBack cb) { this.waitingEntryCallBacks.add(cb); } + public void maybeUpdateCursorBeforeTrimmingConsumedLedger() { + for (ManagedCursor cursor : cursors) { + PositionImpl lastAckedPosition = (PositionImpl) cursor.getMarkDeletedPosition(); + LedgerInfo currPointedLedger = ledgers.get(lastAckedPosition.getLedgerId()); + LedgerInfo nextPointedLedger = Optional.ofNullable(ledgers.higherEntry(lastAckedPosition.getLedgerId())) + .map(Map.Entry::getValue).orElse(null); + + if (currPointedLedger != null) { + if (nextPointedLedger != null) { + if (lastAckedPosition.getEntryId() != -1 && + lastAckedPosition.getEntryId() + 1 >= currPointedLedger.getEntries()) { + lastAckedPosition = new PositionImpl(nextPointedLedger.getLedgerId(), -1); + } + } else { + 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); + } + + if (!lastAckedPosition.equals((PositionImpl) cursor.getMarkDeletedPosition())) { + try { + log.info("Reset cursor:{} to {} since ledger consumed completely", cursor, lastAckedPosition); + lastConfirmedEntry = lastAckedPosition; + updateCursor((ManagedCursorImpl) cursor, lastAckedPosition); + } catch (Exception e) { + log.warn("Failed to reset cursor: {} from {} to {}. Trimming thread will retry next time.", + cursor, cursor.getMarkDeletedPosition(), lastAckedPosition); + log.warn("Caused by", e); + } + } + } + } + private void trimConsumedLedgersInBackground() { trimConsumedLedgersInBackground(Futures.NULL_PROMISE); } @@ -2277,6 +2301,9 @@ void internalTrimConsumedLedgers(CompletableFuture promise) { return; } + // May need to update the cursor position + maybeUpdateCursorBeforeTrimmingConsumedLedger(); + long slowestReaderLedgerId = -1; if (!cursors.hasDurableCursors()) { // At this point the lastLedger will be pointing to the diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java index f6262487b1659..8005693e9f1ea 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java @@ -100,6 +100,7 @@ import org.apache.bookkeeper.mledger.proto.MLDataFormats; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedLedgerInfo; import org.apache.bookkeeper.mledger.proto.MLDataFormats.ManagedLedgerInfo.LedgerInfo; +import org.apache.bookkeeper.mledger.util.Futures; import org.apache.bookkeeper.test.MockedBookKeeperTestCase; import org.apache.commons.lang3.exception.ExceptionUtils; import org.apache.commons.lang3.mutable.MutableObject; @@ -833,7 +834,7 @@ public void testTrimmer() throws Exception { cursor.markDelete(lastPosition); - while (ledger.getNumberOfEntries() != 2) { + while (ledger.getNumberOfEntries() >= 2) { Thread.sleep(10); } } @@ -2906,7 +2907,50 @@ public void testManagedLedgerRollOverIfFull() throws Exception { // all the messages have benn acknowledged // and all the ledgers have been removed except the last ledger Thread.sleep(1000); + ledger.internalTrimConsumedLedgers(Futures.NULL_PROMISE); + Thread.sleep(1000); Assert.assertEquals(ledger.getLedgersInfoAsList().size(), 1); Assert.assertEquals(ledger.getTotalSize(), 0); } + + @Test + public void testExpiredLedgerDeletionAfterManagedLedgerRestart() throws Exception { + ManagedLedgerConfig config = new ManagedLedgerConfig(); + config.setRetentionTime(1, TimeUnit.SECONDS); + config.setMaxEntriesPerLedger(2); + config.setMinimumRolloverTime(1, TimeUnit.MILLISECONDS); + config.setMaximumRolloverTime(500, TimeUnit.MILLISECONDS); + + ManagedLedgerImpl managedLedger = (ManagedLedgerImpl) factory.open("ml_restart_ledger", config); + ManagedCursor cursor = managedLedger.openCursor("c1"); + + for (int i = 0; i < 3; i++) { + managedLedger.addEntry(new byte[1024 * 1024]); + } + + // we have 2 ledgers at the beginning [{entries=2}, {entries=1}] + Assert.assertEquals(managedLedger.getLedgersInfoAsList().size(), 2); + List entries = cursor.readEntries(3); + + for (Entry entry : entries) { + cursor.markDelete(entry.getPosition()); + } + entries.forEach(e -> e.release()); + + // managed-ledger restart + managedLedger.close(); + managedLedger = (ManagedLedgerImpl) factory.open("ml_restart_ledger", config); + + // then we have one more empty ledger after managed-ledger initialization + // the oldest ledger({entries=2}) will be removed when ledger closed + // and now we have [{entries=1}, {entries=0}] + Assert.assertEquals(managedLedger.getLedgersInfoAsList().size(), 2); + + // Now we update the cursors that are still subscribing to ledgers that has been consumed completely + managedLedger.internalTrimConsumedLedgers(Futures.NULL_PROMISE); + + // We only have one empty ledger at last [{entries=0}] + Assert.assertEquals(managedLedger.getLedgersInfoAsList().size(), 1); + Assert.assertEquals(managedLedger.getTotalSize(), 0); + } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/CurrentLedgerRolloverIfFullTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/CurrentLedgerRolloverIfFullTest.java index 1bb8dcb24e4d8..2d5032d59a46c 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/CurrentLedgerRolloverIfFullTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/CurrentLedgerRolloverIfFullTest.java @@ -23,6 +23,7 @@ import lombok.Cleanup; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; +import org.apache.bookkeeper.mledger.util.Futures; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.Message; @@ -99,6 +100,8 @@ public void testCurrentLedgerRolloverIfFull() throws Exception { // trigger a ledger rollover managedLedger.rollCurrentLedgerIfFull(); + Thread.sleep(1000); + managedLedger.trimConsumedLedgersInBackground(Futures.NULL_PROMISE); // the last ledger will be closed and removed and we have one ledger for empty Awaitility.await() From bf63ac010c356eba9216c65b4fd4c5cbfcf6d564 Mon Sep 17 00:00:00 2001 From: wuzhanpeng Date: Wed, 31 Mar 2021 17:53:54 +0800 Subject: [PATCH 2/5] bugfix for cleanup expired data --- .../mledger/impl/ManagedLedgerFactoryImpl.java | 3 +++ .../bookkeeper/mledger/impl/ManagedLedgerImpl.java | 10 ++++++---- .../bookkeeper/mledger/impl/ManagedLedgerTest.java | 11 +++++------ .../service/CurrentLedgerRolloverIfFullTest.java | 2 -- 4 files changed, 14 insertions(+), 12 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java index bed86a140483d..eb8f987c5140c 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerFactoryImpl.java @@ -395,6 +395,9 @@ public void initializeComplete() { log.info("[{}] Successfully initialize managed ledger", name); pendingInitializeLedgers.remove(name, pendingLedger); future.complete(newledger); + + // May need to update the cursor position + newledger.maybeUpdateCursorBeforeTrimmingConsumedLedger(); } @Override 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 08f4420213244..b1c9bf42f9e90 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 @@ -1405,6 +1405,9 @@ public void operationComplete(Void v, Stat stat) { mbean.addLedgerSwitchLatencySample(System.currentTimeMillis() - lastLedgerCreationInitiationTimestamp, TimeUnit.MILLISECONDS); } + + // May need to update the cursor position + maybeUpdateCursorBeforeTrimmingConsumedLedger(); } @Override @@ -2153,7 +2156,9 @@ public void maybeUpdateCursorBeforeTrimmingConsumedLedger() { if (!lastAckedPosition.equals((PositionImpl) cursor.getMarkDeletedPosition())) { try { log.info("Reset cursor:{} to {} since ledger consumed completely", cursor, lastAckedPosition); - lastConfirmedEntry = lastAckedPosition; + if (lastConfirmedEntry.compareTo(lastAckedPosition) < 0) { + lastConfirmedEntry = lastAckedPosition; + } updateCursor((ManagedCursorImpl) cursor, lastAckedPosition); } catch (Exception e) { log.warn("Failed to reset cursor: {} from {} to {}. Trimming thread will retry next time.", @@ -2301,9 +2306,6 @@ void internalTrimConsumedLedgers(CompletableFuture promise) { return; } - // May need to update the cursor position - maybeUpdateCursorBeforeTrimmingConsumedLedger(); - long slowestReaderLedgerId = -1; if (!cursors.hasDurableCursors()) { // At this point the lastLedger will be pointing to the diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java index 8005693e9f1ea..02fb1318a653f 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java @@ -834,7 +834,7 @@ public void testTrimmer() throws Exception { cursor.markDelete(lastPosition); - while (ledger.getNumberOfEntries() >= 2) { + while (ledger.getNumberOfEntries() != 2) { Thread.sleep(10); } } @@ -2907,8 +2907,6 @@ public void testManagedLedgerRollOverIfFull() throws Exception { // all the messages have benn acknowledged // and all the ledgers have been removed except the last ledger Thread.sleep(1000); - ledger.internalTrimConsumedLedgers(Futures.NULL_PROMISE); - Thread.sleep(1000); Assert.assertEquals(ledger.getLedgersInfoAsList().size(), 1); Assert.assertEquals(ledger.getTotalSize(), 0); } @@ -2942,12 +2940,13 @@ public void testExpiredLedgerDeletionAfterManagedLedgerRestart() throws Exceptio managedLedger = (ManagedLedgerImpl) factory.open("ml_restart_ledger", config); // then we have one more empty ledger after managed-ledger initialization - // the oldest ledger({entries=2}) will be removed when ledger closed - // and now we have [{entries=1}, {entries=0}] - Assert.assertEquals(managedLedger.getLedgersInfoAsList().size(), 2); + // and now ledgers are [{entries=2}, {entries=1}, {entries=0}] + Assert.assertTrue(managedLedger.getLedgersInfoAsList().size() >= 2); // Now we update the cursors that are still subscribing to ledgers that has been consumed completely + managedLedger.maybeUpdateCursorBeforeTrimmingConsumedLedger(); managedLedger.internalTrimConsumedLedgers(Futures.NULL_PROMISE); + Thread.sleep(100); // We only have one empty ledger at last [{entries=0}] Assert.assertEquals(managedLedger.getLedgersInfoAsList().size(), 1); diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/CurrentLedgerRolloverIfFullTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/CurrentLedgerRolloverIfFullTest.java index 2d5032d59a46c..08cd1a69f67b4 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/CurrentLedgerRolloverIfFullTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/CurrentLedgerRolloverIfFullTest.java @@ -100,8 +100,6 @@ public void testCurrentLedgerRolloverIfFull() throws Exception { // trigger a ledger rollover managedLedger.rollCurrentLedgerIfFull(); - Thread.sleep(1000); - managedLedger.trimConsumedLedgersInBackground(Futures.NULL_PROMISE); // the last ledger will be closed and removed and we have one ledger for empty Awaitility.await() From 9151921ce19f568f4ec1f2f0dfb084ebe6d1e8f2 Mon Sep 17 00:00:00 2001 From: wuzhanpeng Date: Wed, 31 Mar 2021 17:58:13 +0800 Subject: [PATCH 3/5] remove unused imports --- .../pulsar/broker/service/CurrentLedgerRolloverIfFullTest.java | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/CurrentLedgerRolloverIfFullTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/CurrentLedgerRolloverIfFullTest.java index 08cd1a69f67b4..1bb8dcb24e4d8 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/CurrentLedgerRolloverIfFullTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/CurrentLedgerRolloverIfFullTest.java @@ -23,7 +23,6 @@ import lombok.Cleanup; import org.apache.bookkeeper.mledger.ManagedLedgerConfig; import org.apache.bookkeeper.mledger.impl.ManagedLedgerImpl; -import org.apache.bookkeeper.mledger.util.Futures; import org.apache.pulsar.broker.service.persistent.PersistentTopic; import org.apache.pulsar.client.api.Consumer; import org.apache.pulsar.client.api.Message; From 35f2a2aaa590d96fa6290337b8e43e6f01280a52 Mon Sep 17 00:00:00 2001 From: wuzhanpeng Date: Tue, 6 Apr 2021 14:12:10 +0800 Subject: [PATCH 4/5] avoid md-entry and last-confirm-entry being equal --- .../bookkeeper/mledger/impl/ManagedCursorPropertiesTest.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorPropertiesTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorPropertiesTest.java index 43720fd6dd0e9..eb682644718cd 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorPropertiesTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/ManagedCursorPropertiesTest.java @@ -45,6 +45,7 @@ void testPropertiesClose() throws Exception { ledger.addEntry("entry-1".getBytes()); ledger.addEntry("entry-2".getBytes()); Position p3 = ledger.addEntry("entry-3".getBytes()); + ledger.addEntry("entry-4".getBytes()); Map properties = new TreeMap<>(); properties.put("a", 1L); @@ -82,6 +83,7 @@ void testPropertiesRecoveryAfterCrash() throws Exception { ledger.addEntry("entry-1".getBytes()); ledger.addEntry("entry-2".getBytes()); Position p3 = ledger.addEntry("entry-3".getBytes()); + ledger.addEntry("entry-4".getBytes()); Map properties = new TreeMap<>(); properties.put("a", 1L); @@ -113,6 +115,7 @@ void testPropertiesOnDelete() throws Exception { ledger.addEntry("entry-1".getBytes()); Position p2 = ledger.addEntry("entry-2".getBytes()); Position p3 = ledger.addEntry("entry-3".getBytes()); + ledger.addEntry("entry-4".getBytes()); Map properties = new TreeMap<>(); properties.put("a", 1L); From fef56bd0f7c8938b9532788826b8479a7ef2da8a Mon Sep 17 00:00:00 2001 From: wuzhanpeng Date: Fri, 9 Apr 2021 10:27:15 +0800 Subject: [PATCH 5/5] avoid modifying lastConfirmedEntry --- .../mledger/impl/ManagedCursorImpl.java | 27 ++++++++++++++----- .../mledger/impl/ManagedLedgerImpl.java | 3 --- 2 files changed, 21 insertions(+), 9 deletions(-) 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 81e8399777704..13000f08b6ae7 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 @@ -1631,13 +1631,28 @@ public void asyncMarkDelete(final Position position, Map propertie } if (((PositionImpl) ledger.getLastConfirmedEntry()).compareTo(newPosition) < 0) { - if (log.isDebugEnabled()) { - log.debug( - "[{}] Failed mark delete due to invalid markDelete {} is ahead of last-confirmed-entry {} for cursor [{}]", - ledger.getName(), position, ledger.getLastConfirmedEntry(), name); + boolean shouldCursorMoveForward = false; + try { + long ledgerEntries = ledger.getLedgerInfo(markDeletePosition.getLedgerId()).get().getEntries(); + Long nextValidLedger = ledger.getNextValidLedger(ledger.getLastConfirmedEntry().getLedgerId()); + shouldCursorMoveForward = (markDeletePosition.getEntryId() + 1 >= ledgerEntries) + && (newPosition.getLedgerId() == nextValidLedger); + } catch (Exception e) { + log.warn("Failed to get ledger entries while setting mark-delete-position", e); + } + + if (shouldCursorMoveForward) { + log.info("[{}] move mark-delete-position from {} to {} since all the entries have been consumed", + ledger.getName(), markDeletePosition, newPosition); + } else { + if (log.isDebugEnabled()) { + log.debug( + "[{}] Failed mark delete due to invalid markDelete {} is ahead of last-confirmed-entry {} for cursor [{}]", + ledger.getName(), position, ledger.getLastConfirmedEntry(), name); + } + callback.markDeleteFailed(new ManagedLedgerException("Invalid mark deleted position"), ctx); + return; } - callback.markDeleteFailed(new ManagedLedgerException("Invalid mark deleted position"), ctx); - return; } lock.writeLock().lock(); 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 b1c9bf42f9e90..5328d07f8c154 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 @@ -2156,9 +2156,6 @@ public void maybeUpdateCursorBeforeTrimmingConsumedLedger() { if (!lastAckedPosition.equals((PositionImpl) cursor.getMarkDeletedPosition())) { try { log.info("Reset cursor:{} to {} since ledger consumed completely", cursor, lastAckedPosition); - if (lastConfirmedEntry.compareTo(lastAckedPosition) < 0) { - lastConfirmedEntry = lastAckedPosition; - } updateCursor((ManagedCursorImpl) cursor, lastAckedPosition); } catch (Exception e) { log.warn("Failed to reset cursor: {} from {} to {}. Trimming thread will retry next time.",