From 96f7eb2a54385c3dd5e03ebd7e5e4f9b575e6b8b Mon Sep 17 00:00:00 2001 From: hangc0276 Date: Mon, 14 Dec 2020 09:10:45 +0800 Subject: [PATCH 1/6] mmake ledger rollover check task internal --- .../bookkeeper/mledger/ManagedLedger.java | 1 + .../mledger/impl/ManagedLedgerImpl.java | 16 +++++++-- .../mledger/impl/ManagedLedgerTest.java | 33 +++++++++++++++++++ .../pulsar/broker/service/BrokerService.java | 23 ------------- 4 files changed, 48 insertions(+), 25 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 f75a63903ff16..4f4c226a91e2a 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 @@ -583,6 +583,7 @@ void asyncSetProperties(Map properties, final AsyncCallbacks.Upd /** * Roll current ledger if it is full */ + @Deprecated void rollCurrentLedgerIfFull(); /** 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 4a49406bf4cb2..0880fb884645c 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 @@ -176,6 +176,7 @@ public class ManagedLedgerImpl implements ManagedLedger, CreateCallback { final EntryCache entryCache; private ScheduledFuture timeoutTask; + private ScheduledFuture checkLedgerRollTask; /** * This lock is held while the ledgers list or propertiesMap is updated asynchronously on the metadata store. Since we use the store @@ -1318,6 +1319,10 @@ public synchronized void asyncClose(final CloseCallback callback, final Object c this.timeoutTask.cancel(false); } + if (this.checkLedgerRollTask != null) { + this.checkLedgerRollTask.cancel(false); + } + } private void closeAllCursors(CloseCallback callback, final Object ctx) { @@ -1548,7 +1553,7 @@ synchronized void createLedgerAfterClosed() { asyncCreateLedger(bookKeeper, config, digestType, this, Collections.emptyMap()); } - @Override + @VisibleForTesting public void rollCurrentLedgerIfFull() { log.info("[{}] Start checking if current ledger is full", name); if (currentLedgerEntries > 0 && currentLedgerIsFull()) { @@ -1561,7 +1566,7 @@ public void closeComplete(int rc, LedgerHandle lh, Object o) { lh.getId()); if (rc == BKException.Code.OK) { - log.debug("Successfuly closed ledger {}", lh.getId()); + log.debug("Successfully closed ledger {}", lh.getId()); } else { log.warn("Error when closing ledger {}. Status={}", lh.getId(), BKException.getMessage(rc)); } @@ -3465,6 +3470,13 @@ private void scheduleTimeoutTask() { checkReadTimeout(); }), timeoutSec, timeoutSec, TimeUnit.SECONDS); } + + if (config.getMaximumRolloverTimeMs() > 0) { + long interval = config.getMaximumRolloverTimeMs(); + this.checkLedgerRollTask = this.scheduledExecutor.scheduleAtFixedRate(safeRun(() -> { + rollCurrentLedgerIfFull(); + }), interval, interval, TimeUnit.MILLISECONDS); + } } private void checkAddTimeout() { 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 ff73f38ef0919..80744d9edd3f7 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 @@ -2825,4 +2825,37 @@ public static void retryStrategically(Predicate predicate, int retryCount, Thread.sleep(intSleepTimeInMillis + (intSleepTimeInMillis * i)); } } + + @Test + public void testManagedLedgerRollOverIfFull() 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 ledger = (ManagedLedgerImpl)factory.open("test_managedLedger_rollOver", config); + ManagedCursor cursor = ledger.openCursor("c1"); + + int msgNum = 10; + + for (int i = 0; i < msgNum; i++) { + ledger.addEntry(new byte[1024 * 1024]); + } + + Assert.assertEquals(ledger.getLedgersInfoAsList().size(), msgNum / 2); + List entries = cursor.readEntries(msgNum); + Assert.assertEquals(msgNum, entries.size()); + + for (Entry entry : entries) { + cursor.markDelete(entry.getPosition()); + } + entries.forEach(e -> e.release()); + + // all the messages have benn acknowledged + // and all the ledgers have been removed except the last ledger + Thread.sleep(2000); + Assert.assertEquals(ledger.getLedgersInfoAsList().size(), 1); + Assert.assertEquals(ledger.getCurrentLedgerSize(), 0); + } } 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 2e8fdf71a4e74..8faa4621a235c 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 @@ -216,7 +216,6 @@ public class BrokerService implements Closeable, ZooKeeperCacheListener { - if (t instanceof PersistentTopic) { - Optional.ofNullable(((PersistentTopic) t).getManagedLedger()).ifPresent( - managedLedger -> { - managedLedger.rollCurrentLedgerIfFull(); - } - ); - } - }); - } - public void checkMessageDeduplicationInfo() { forEachTopic(Topic::checkMessageDeduplicationInfo); } From 9465d2903dbf458df1f52d62e0428acb6f0e8c2b Mon Sep 17 00:00:00 2001 From: hangc0276 Date: Mon, 14 Dec 2020 21:15:57 +0800 Subject: [PATCH 2/6] move roll over ledger task into individual method. --- .../org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java | 4 ++++ 1 file changed, 4 insertions(+) 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 0880fb884645c..ba5358a42db37 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 @@ -383,6 +383,8 @@ public void operationFailed(MetaStoreException e) { }); scheduleTimeoutTask(); + + scheduleRollOverLedgerTask(); } private synchronized void initializeBookKeeper(final ManagedLedgerInitializeLedgerCallback callback) { @@ -3470,7 +3472,9 @@ private void scheduleTimeoutTask() { checkReadTimeout(); }), timeoutSec, timeoutSec, TimeUnit.SECONDS); } + } + private void scheduleRollOverLedgerTask() { if (config.getMaximumRolloverTimeMs() > 0) { long interval = config.getMaximumRolloverTimeMs(); this.checkLedgerRollTask = this.scheduledExecutor.scheduleAtFixedRate(safeRun(() -> { From facd007dfc6540a981c05a3924381edfd8bca48f Mon Sep 17 00:00:00 2001 From: hangc0276 Date: Thu, 24 Dec 2020 09:33:09 +0800 Subject: [PATCH 3/6] mark rollCurrentLedgerIfFull interface as Deprecated --- .../org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java | 1 + 1 file changed, 1 insertion(+) 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 ba5358a42db37..e295566792e73 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 @@ -1556,6 +1556,7 @@ synchronized void createLedgerAfterClosed() { } @VisibleForTesting + @Override public void rollCurrentLedgerIfFull() { log.info("[{}] Start checking if current ledger is full", name); if (currentLedgerEntries > 0 && currentLedgerIsFull()) { From 534326d168f86e98e60ac9c3c1c106235dbd8391 Mon Sep 17 00:00:00 2001 From: hangc0276 Date: Thu, 7 Jan 2021 13:08:37 +0800 Subject: [PATCH 4/6] fix run test case failed. --- .../org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 80744d9edd3f7..5abd15d3be7ca 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 @@ -2854,7 +2854,7 @@ public void testManagedLedgerRollOverIfFull() throws Exception { // all the messages have benn acknowledged // and all the ledgers have been removed except the last ledger - Thread.sleep(2000); + Thread.sleep(5000); Assert.assertEquals(ledger.getLedgersInfoAsList().size(), 1); Assert.assertEquals(ledger.getCurrentLedgerSize(), 0); } From 766e9c69dcb262dd53d8cdc1b85ffe430d59b53b Mon Sep 17 00:00:00 2001 From: hangc0276 Date: Thu, 7 Jan 2021 20:36:13 +0800 Subject: [PATCH 5/6] fix run test case failed. --- .../apache/bookkeeper/mledger/impl/ManagedLedgerTest.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) 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 5abd15d3be7ca..c14518e515f2d 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 @@ -18,6 +18,7 @@ */ package org.apache.bookkeeper.mledger.impl; +import org.apache.bookkeeper.mledger.ManagedLedgerInfo; import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyInt; @@ -2854,8 +2855,8 @@ public void testManagedLedgerRollOverIfFull() throws Exception { // all the messages have benn acknowledged // and all the ledgers have been removed except the last ledger - Thread.sleep(5000); - Assert.assertEquals(ledger.getLedgersInfoAsList().size(), 1); + Thread.sleep(1000); + Assert.assertEquals(ledger.getLedgersInfoAsList().size(), 2); Assert.assertEquals(ledger.getCurrentLedgerSize(), 0); } } From f94c36e311b6215ea4f3eec54e7b00e1cf198890 Mon Sep 17 00:00:00 2001 From: hangc0276 Date: Thu, 7 Jan 2021 20:52:03 +0800 Subject: [PATCH 6/6] fix run test case failed. --- .../org/apache/bookkeeper/mledger/impl/ManagedLedgerTest.java | 1 - 1 file changed, 1 deletion(-) 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 c14518e515f2d..11cb57f7c2d86 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 @@ -18,7 +18,6 @@ */ package org.apache.bookkeeper.mledger.impl; -import org.apache.bookkeeper.mledger.ManagedLedgerInfo; import org.apache.pulsar.common.policies.data.EnsemblePlacementPolicyConfig; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyInt;