From 6d729ffd3458ee234eb18447fa07e0750c109578 Mon Sep 17 00:00:00 2001 From: penghui Date: Tue, 20 Jul 2021 17:58:58 +0800 Subject: [PATCH 1/3] Invalidate the read handle after all cursors consumed. Currently, the read ReadHandle only invalidate when removing the ledger from the ManagedLedger. If the ManagedLedger have many ledgers(the topic might retain infinite data), we will get oom on the direct memory since we are using 1MB read cache by default for the offloaded data ReadHandle. If all the cursors are consumed the data, the ReadHandle can be closed safety. And if a cursor reset to an earlier position to consume the historical data, the ReadHandle will be reopen again. --- .../mledger/impl/ManagedLedgerImpl.java | 4 +- .../mledger/impl/ManagedLedgerTest.java | 61 +++++++++++++++++++ 2 files changed, 63 insertions(+), 2 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 4981910ce9b84..d98c0905b9b2e 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 @@ -151,7 +151,7 @@ public class ManagedLedgerImpl implements ManagedLedger, CreateCallback { protected Map propertiesMap; protected final MetaStore store; - private final ConcurrentLongHashMap> ledgerCache = new ConcurrentLongHashMap<>( + protected final ConcurrentLongHashMap> ledgerCache = new ConcurrentLongHashMap<>( 16 /* initial capacity */, 1 /* number of sections */); protected final NavigableMap ledgers = new ConcurrentSkipListMap<>(); private volatile Stat ledgersStat; @@ -2410,7 +2410,7 @@ void internalTrimLedgers(boolean isTruncate, CompletableFuture promise) { if (log.isDebugEnabled()) { log.debug("[{}] Ledger {} not deleted. Neither expired nor over-quota", name, ls.getLedgerId()); } - break; + invalidateReadHandle(ls.getLedgerId()); } } 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 25c1725bdb4cc..3b66149d6d084 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 @@ -3097,4 +3097,65 @@ public void testOpEntryAdd_toString_doesNotThrowNPE(){ ", dataLength=" + dataLength + '}'; } + + @Test + public void testInvalidateReadHandleWhenDeleteLedger() throws Exception { + ManagedLedgerConfig config = new ManagedLedgerConfig(); + config.setMaxEntriesPerLedger(1); + + // Verify the read handle should be invalidated after ledger been removed. + ManagedLedgerImpl ledger = (ManagedLedgerImpl)factory.open("testInvalidateReadHandleWhenDeleteLedger", config); + ManagedCursor cursor = ledger.openCursor("test-cursor"); + ManagedCursor cursor2 = ledger.openCursor("test-cursor2"); + final int entries = 3; + for (int i = 0; i < entries; i++) { + ledger.addEntry(String.valueOf(i).getBytes(Encoding)); + } + List entryList = cursor.readEntries(3); + assertEquals(entryList.size(), 3); + assertEquals(ledger.ledgers.size(), 3); + assertEquals(ledger.ledgerCache.size(), 2); + cursor.clearBacklog(); + cursor2.clearBacklog(); + ledger.trimConsumedLedgersInBackground(Futures.NULL_PROMISE); + Awaitility.await().untilAsserted(() -> { + assertEquals(ledger.ledgers.size(), 1); + assertEquals(ledger.ledgerCache.size(), 0); + }); + + cursor.close(); + cursor2.close(); + ledger.close(); + } + + @Test + public void testInvalidateReadHandleWhenConsumed() throws Exception { + ManagedLedgerConfig config = new ManagedLedgerConfig(); + config.setMaxEntriesPerLedger(1); + // Verify the read handle should be invalidated when all cursors consumed + // even if the ledger can not been removed due to the data retention + config.setRetentionSizeInMB(50); + config.setRetentionTime(1, TimeUnit.DAYS); + ManagedLedgerImpl ledger = (ManagedLedgerImpl)factory.open("testInvalidateReadHandleWhenConsumed", config); + ManagedCursor cursor = ledger.openCursor("test-cursor"); + ManagedCursor cursor2 = ledger.openCursor("test-cursor2"); + final int entries = 3; + for (int i = 0; i < entries; i++) { + ledger.addEntry(String.valueOf(i).getBytes(Encoding)); + } + List entryList = cursor.readEntries(3); + assertEquals(entryList.size(), 3); + assertEquals(ledger.ledgers.size(), 3); + assertEquals(ledger.ledgerCache.size(), 2); + cursor.clearBacklog(); + cursor2.clearBacklog(); + ledger.trimConsumedLedgersInBackground(Futures.NULL_PROMISE); + Awaitility.await().untilAsserted(() -> { + assertEquals(ledger.ledgers.size(), 3); + assertEquals(ledger.ledgerCache.size(), 0); + }); + cursor.close(); + cursor2.close(); + ledger.close(); + } } From a497146f40c8fc8d5c74fbe03a0f13149b03000e Mon Sep 17 00:00:00 2001 From: penghui Date: Tue, 20 Jul 2021 19:02:09 +0800 Subject: [PATCH 2/3] Add tests. --- .../mledger/impl/ManagedLedgerTest.java | 15 +++++++++++++++ 1 file changed, 15 insertions(+) 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 3b66149d6d084..cd805379ac059 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 @@ -3154,8 +3154,23 @@ public void testInvalidateReadHandleWhenConsumed() throws Exception { assertEquals(ledger.ledgers.size(), 3); assertEquals(ledger.ledgerCache.size(), 0); }); + + // Verify the ReadHandle can be reopened. + ManagedCursor cursor3 = ledger.openCursor("test-cursor3", InitialPosition.Earliest); + entryList = cursor3.readEntries(3); + assertEquals(entryList.size(), 3); + assertEquals(ledger.ledgerCache.size(), 2); + cursor3.clearBacklog(); + ledger.trimConsumedLedgersInBackground(Futures.NULL_PROMISE); + Awaitility.await().untilAsserted(() -> { + assertEquals(ledger.ledgers.size(), 3); + assertEquals(ledger.ledgerCache.size(), 0); + }); + + cursor.close(); cursor2.close(); + cursor3.close(); ledger.close(); } } From 33bb631afec637d15c53cdfedd436fa42c822232 Mon Sep 17 00:00:00 2001 From: penghui Date: Tue, 20 Jul 2021 23:04:13 +0800 Subject: [PATCH 3/3] Apply comment. --- .../org/apache/bookkeeper/mledger/impl/ManagedLedgerImpl.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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 d98c0905b9b2e..e611d634f9fb4 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 @@ -151,7 +151,7 @@ public class ManagedLedgerImpl implements ManagedLedger, CreateCallback { protected Map propertiesMap; protected final MetaStore store; - protected final ConcurrentLongHashMap> ledgerCache = new ConcurrentLongHashMap<>( + final ConcurrentLongHashMap> ledgerCache = new ConcurrentLongHashMap<>( 16 /* initial capacity */, 1 /* number of sections */); protected final NavigableMap ledgers = new ConcurrentSkipListMap<>(); private volatile Stat ledgersStat;