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..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 @@ -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 @@ -382,6 +383,8 @@ public void operationFailed(MetaStoreException e) { }); scheduleTimeoutTask(); + + scheduleRollOverLedgerTask(); } private synchronized void initializeBookKeeper(final ManagedLedgerInitializeLedgerCallback callback) { @@ -1318,6 +1321,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,6 +1555,7 @@ synchronized void createLedgerAfterClosed() { asyncCreateLedger(bookKeeper, config, digestType, this, Collections.emptyMap()); } + @VisibleForTesting @Override public void rollCurrentLedgerIfFull() { log.info("[{}] Start checking if current ledger is full", name); @@ -1561,7 +1569,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)); } @@ -3467,6 +3475,15 @@ private void scheduleTimeoutTask() { } } + private void scheduleRollOverLedgerTask() { + if (config.getMaximumRolloverTimeMs() > 0) { + long interval = config.getMaximumRolloverTimeMs(); + this.checkLedgerRollTask = this.scheduledExecutor.scheduleAtFixedRate(safeRun(() -> { + rollCurrentLedgerIfFull(); + }), interval, interval, TimeUnit.MILLISECONDS); + } + } + private void checkAddTimeout() { long timeoutSec = config.getAddEntryTimeoutSeconds(); if (timeoutSec < 1) { 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..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 @@ -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(1000); + Assert.assertEquals(ledger.getLedgersInfoAsList().size(), 2); + 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); }