From f71b044756ba053678cd7b8e229377566e98dfdc Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ra=C3=BAl=20Gracia?= Date: Fri, 27 Aug 2021 19:07:22 +0200 Subject: [PATCH 1/5] Fixed entry log GC for entrylogPerLedgerEnabled. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Raúl Gracia --- .../apache/bookkeeper/bookie/EntryLogger.java | 30 +++++++++++++++++++ .../bookie/GarbageCollectorThread.java | 18 +++++++++-- 2 files changed, 45 insertions(+), 3 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/EntryLogger.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/EntryLogger.java index 504adfa4214..49a9ca40594 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/EntryLogger.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/EntryLogger.java @@ -482,6 +482,27 @@ long getLeastUnflushedLogId() { return recentlyCreatedEntryLogsStatus.getLeastUnflushedLogId(); } + /** + * Get the last log id created so far. If entryLogPerLedger is enabled, the Garbage Collector + * process needs to look beyond the least unflushed entry log file, as there may be entry logs + * ready to be garbage collected. + * + * @return last entry log id created. + */ + long getLastLogId() { + return recentlyCreatedEntryLogsStatus.getLastLogId(); + } + + /** + * Returns whether the current log id exists and has been rotated already. + * + * @param entryLogId EntryLog id to check. + * @return Whether the given entryLogId exists and has been rotated. + */ + boolean isFlushedEntryLog(Long entryLogId) { + return recentlyCreatedEntryLogsStatus.isFlushedEntryLog(entryLogId); + } + long getPreviousAllocatedEntryLogId() { return entryLoggerAllocator.getPreallocatedLogId(); } @@ -1249,5 +1270,14 @@ synchronized void flushRotatedEntryLog(Long entryLogId) { synchronized long getLeastUnflushedLogId() { return leastUnflushedLogId; } + + synchronized long getLastLogId() { + return !entryLogsStatusMap.isEmpty() ? entryLogsStatusMap.lastKey() : 0; + } + + synchronized boolean isFlushedEntryLog(Long entryLogId) { + return entryLogsStatusMap.containsKey(entryLogId) && entryLogsStatusMap.get(entryLogId) + || entryLogId < leastUnflushedLogId; + } } } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/GarbageCollectorThread.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/GarbageCollectorThread.java index fb548905ff7..325aff8b2ee 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/GarbageCollectorThread.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/GarbageCollectorThread.java @@ -37,6 +37,7 @@ import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicLong; +import java.util.function.Supplier; import lombok.Getter; import org.apache.bookkeeper.bookie.GarbageCollector.GarbageCleaner; @@ -586,9 +587,11 @@ protected Map extractMetaFromEntryLogs(Map finalEntryLog = () -> conf.isEntryLogPerLedgerEnabled() ? entryLogger.getLastLogId() : + entryLogger.getLeastUnflushedLogId(); boolean hasExceptionWhenScan = false; - for (long entryLogId = scannedLogId; entryLogId < curLogId; entryLogId++) { + boolean increaseScannedLogId = true; + for (long entryLogId = scannedLogId; entryLogId < finalEntryLog.get(); entryLogId++) { // Comb the current entry log file if it has not already been extracted. if (entryLogMetaMap.containsKey(entryLogId)) { continue; @@ -600,6 +603,15 @@ protected Map extractMetaFromEntryLogs(Map extractMetaFromEntryLogs(Map Date: Sat, 28 Aug 2021 17:52:03 +0200 Subject: [PATCH 2/5] Added test and comments. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Raúl --- .../bookie/GarbageCollectorThread.java | 10 ++-- .../bookkeeper/bookie/CompactionTest.java | 53 +++++++++++++++++++ 2 files changed, 59 insertions(+), 4 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/GarbageCollectorThread.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/GarbageCollectorThread.java index 325aff8b2ee..cafbf53eef5 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/GarbageCollectorThread.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/GarbageCollectorThread.java @@ -584,9 +584,10 @@ protected void compactEntryLog(EntryLogMetadata entryLogMeta) { * @throws IOException */ protected Map extractMetaFromEntryLogs(Map entryLogMetaMap) { - // Extract it for every entry log except for the current one. - // Entry Log ID's are just a long value that starts at 0 and increments - // by 1 when the log fills up and we roll to a new one. + // Entry Log ID's are just a long value that starts at 0 and increments by 1 when the log fills up and we roll + // to a new one. We scan entry logs as follows: + // - entryLogPerLedgerEnabled is false: Extract it for every entry log except for the current one (un-flushed). + // - entryLogPerLedgerEnabled is true: Scan all flushed entry logs up to the highest known id. Supplier finalEntryLog = () -> conf.isEntryLogPerLedgerEnabled() ? entryLogger.getLastLogId() : entryLogger.getLeastUnflushedLogId(); boolean hasExceptionWhenScan = false; @@ -631,7 +632,8 @@ protected Map extractMetaFromEntryLogs(Map 0); } + @Test + public void testMinorCompactionWithEntryLogPerLedgerEnabled() throws Exception { + // restart bookies + restartBookies(c-> { + c.setMajorCompactionThreshold(0.0f); + c.setGcWaitTime(60000); + c.setMinorCompactionInterval(120000); + c.setMajorCompactionInterval(240000); + c.setForceAllowCompaction(true); + c.setEntryLogPerLedgerEnabled(true); + return c; + }); + + // prepare data + LedgerHandle[] lhs = prepareData(3, false); + + for (LedgerHandle lh : lhs) { + lh.close(); + } + + long lastMinorCompactionTime = getGCThread().lastMinorCompactionTime; + long lastMajorCompactionTime = getGCThread().lastMajorCompactionTime; + assertFalse(getGCThread().enableMajorCompaction); + assertTrue(getGCThread().enableMinorCompaction); + + // remove ledgers 1 and 2 + bkc.deleteLedger(lhs[1].getId()); + bkc.deleteLedger(lhs[2].getId()); + + LOG.info("Finished deleting the ledgers contains most entries."); + getGCThread().triggerGC(true, false, false).get(); + + assertEquals(lastMajorCompactionTime, getGCThread().lastMajorCompactionTime); + assertTrue(getGCThread().lastMinorCompactionTime > lastMinorCompactionTime); + + // At this point, we have the following state of ledgers end entry logs: + // L0 (not deleted) -> E0 (un-flushed): Entry log should exist. + // L1 (deleted) -> E1 (un-flushed): Entry log should exist as un-flushed entry logs are not considered for GC. + // L2 (deleted) -> E2 (flushed): Entry log should have been garbage collected. + // E3 (flushed): Entry log should have been garbage collected. + // E4 (un-flushed): Entry log should exist as un-flushed entry logs are not considered for GC. + assertTrue("Entry log file 0.log is not available, which is not expected " + tmpDirs.get(0), + TestUtils.hasLogFiles(tmpDirs.get(0), false, 0)); + assertTrue("Entry log file 1.log is not available, which is not expected " + tmpDirs.get(0), + TestUtils.hasLogFiles(tmpDirs.get(0), false, 1)); + assertTrue("Entry log file 4.log is not available, which is not expected " + tmpDirs.get(0), + TestUtils.hasLogFiles(tmpDirs.get(0), false, 4)); + assertFalse("Entry log file 2.log is available, which is not expected" + tmpDirs.get(0), + TestUtils.hasLogFiles(tmpDirs.get(0), false, 2)); + assertFalse("Entry log file 3.log is available, which is not expected" + tmpDirs.get(0), + TestUtils.hasLogFiles(tmpDirs.get(0), false, 3)); + } + @Test public void testMinorCompactionWithNoWritableLedgerDirs() throws Exception { // prepare data From b036f73d9bcd420d16ea7dc727425ae1a307cc71 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ra=C3=BAl?= Date: Mon, 30 Aug 2021 11:14:39 +0200 Subject: [PATCH 3/5] Working on tests. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Raúl --- .../org/apache/bookkeeper/bookie/CompactionTest.java | 9 +++++---- .../test/java/org/apache/bookkeeper/util/TestUtils.java | 2 +- 2 files changed, 6 insertions(+), 5 deletions(-) diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/CompactionTest.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/CompactionTest.java index 265bb44b4d6..2695ef1d475 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/CompactionTest.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/CompactionTest.java @@ -587,14 +587,15 @@ public void testMinorCompactionWithEntryLogPerLedgerEnabled() throws Exception { // E4 (un-flushed): Entry log should exist as un-flushed entry logs are not considered for GC. assertTrue("Entry log file 0.log is not available, which is not expected " + tmpDirs.get(0), TestUtils.hasLogFiles(tmpDirs.get(0), false, 0)); + verifyLedger(lhs[0].getId(), 0, lhs[0].getLastAddConfirmed()); assertTrue("Entry log file 1.log is not available, which is not expected " + tmpDirs.get(0), TestUtils.hasLogFiles(tmpDirs.get(0), false, 1)); assertTrue("Entry log file 4.log is not available, which is not expected " + tmpDirs.get(0), TestUtils.hasLogFiles(tmpDirs.get(0), false, 4)); - assertFalse("Entry log file 2.log is available, which is not expected" + tmpDirs.get(0), - TestUtils.hasLogFiles(tmpDirs.get(0), false, 2)); - assertFalse("Entry log file 3.log is available, which is not expected" + tmpDirs.get(0), - TestUtils.hasLogFiles(tmpDirs.get(0), false, 3)); + for (File ledgerDirectory : tmpDirs) { + assertFalse("Found entry log files [2, 3].log that should have been compacted in ledgerDirectory: " + + ledgerDirectory, TestUtils.hasLogFiles(ledgerDirectory, false, 2, 3)); + } } @Test diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/util/TestUtils.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/util/TestUtils.java index 462d4729bd4..7281f8f0dc9 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/util/TestUtils.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/util/TestUtils.java @@ -50,7 +50,7 @@ public static String buildStatsCounterPathFromBookieID(BookieId bookieId) { } public static boolean hasLogFiles(File ledgerDirectory, boolean partial, Integer... logsId) { - boolean result = partial ? false : true; + boolean result = !partial; Set logs = new HashSet(); for (File file : BookieImpl.getCurrentDirectory(ledgerDirectory).listFiles()) { if (file.isFile()) { From b6a3ea9be599f6f4f582c7c40020b1ad926931d7 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ra=C3=BAl?= Date: Mon, 30 Aug 2021 11:48:01 +0200 Subject: [PATCH 4/5] Improving test readability. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Raúl --- .../bookkeeper/bookie/CompactionTest.java | 15 +++----- .../org/apache/bookkeeper/util/TestUtils.java | 34 ++++++++++++++----- 2 files changed, 29 insertions(+), 20 deletions(-) diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/CompactionTest.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/CompactionTest.java index 2695ef1d475..beac55d6cc3 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/CompactionTest.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/CompactionTest.java @@ -585,17 +585,10 @@ public void testMinorCompactionWithEntryLogPerLedgerEnabled() throws Exception { // L2 (deleted) -> E2 (flushed): Entry log should have been garbage collected. // E3 (flushed): Entry log should have been garbage collected. // E4 (un-flushed): Entry log should exist as un-flushed entry logs are not considered for GC. - assertTrue("Entry log file 0.log is not available, which is not expected " + tmpDirs.get(0), - TestUtils.hasLogFiles(tmpDirs.get(0), false, 0)); - verifyLedger(lhs[0].getId(), 0, lhs[0].getLastAddConfirmed()); - assertTrue("Entry log file 1.log is not available, which is not expected " + tmpDirs.get(0), - TestUtils.hasLogFiles(tmpDirs.get(0), false, 1)); - assertTrue("Entry log file 4.log is not available, which is not expected " + tmpDirs.get(0), - TestUtils.hasLogFiles(tmpDirs.get(0), false, 4)); - for (File ledgerDirectory : tmpDirs) { - assertFalse("Found entry log files [2, 3].log that should have been compacted in ledgerDirectory: " - + ledgerDirectory, TestUtils.hasLogFiles(ledgerDirectory, false, 2, 3)); - } + assertTrue("Not found entry log files [0, 1, 4].log that should not have been compacted in: " + + tmpDirs.get(0), TestUtils.hasAllLogFiles(tmpDirs.get(0), 0, 1, 4)); + assertTrue("Found entry log files [2, 3].log that should have been compacted in ledgerDirectory: " + + tmpDirs.get(0), TestUtils.hasNoneLogFiles(tmpDirs.get(0), 2, 3)); } @Test diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/util/TestUtils.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/util/TestUtils.java index 7281f8f0dc9..27f1abbb96a 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/util/TestUtils.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/util/TestUtils.java @@ -22,6 +22,7 @@ package org.apache.bookkeeper.util; import java.io.File; +import java.util.Arrays; import java.util.Collection; import java.util.HashSet; import java.util.Set; @@ -49,9 +50,31 @@ public static String buildStatsCounterPathFromBookieID(BookieId bookieId) { return bookieId.toString().replace('.', '_').replace('-', '_').replace(":", "_"); } + public static boolean hasAllLogFiles(File ledgerDirectory, Integer... logsId) { + Set logs = findEntryLogFileIds(ledgerDirectory); + return logs.containsAll(Arrays.asList(logsId)); + } + + public static boolean hasNoneLogFiles(File ledgerDirectory, Integer... logsId) { + Set logs = findEntryLogFileIds(ledgerDirectory); + return Arrays.stream(logsId).noneMatch(logs::contains); + } + public static boolean hasLogFiles(File ledgerDirectory, boolean partial, Integer... logsId) { boolean result = !partial; - Set logs = new HashSet(); + Set logs = findEntryLogFileIds(ledgerDirectory); + for (Integer logId : logsId) { + boolean exist = logs.contains(logId); + if ((partial && exist) + || (!partial && !exist)) { + return !result; + } + } + return result; + } + + private static Set findEntryLogFileIds(File ledgerDirectory) { + Set logs = new HashSet<>(); for (File file : BookieImpl.getCurrentDirectory(ledgerDirectory).listFiles()) { if (file.isFile()) { String name = file.getName(); @@ -61,14 +84,7 @@ public static boolean hasLogFiles(File ledgerDirectory, boolean partial, Integer logs.add(Integer.parseInt(name.split("\\.")[0], 16)); } } - for (Integer logId : logsId) { - boolean exist = logs.contains(logId); - if ((partial && exist) - || (!partial && !exist)) { - return !result; - } - } - return result; + return logs; } public static void waitUntilLacUpdated(ReadHandle rh, long newLac) throws Exception { From 97f417b5ba927b21be5d34e00aaad402df0caaf6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ra=C3=BAl?= Date: Mon, 30 Aug 2021 12:35:39 +0200 Subject: [PATCH 5/5] Wait for log to be flushed before assertions. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Signed-off-by: Raúl --- .../bookkeeper/bookie/CompactionTest.java | 23 +++++++++++++++++++ 1 file changed, 23 insertions(+) diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/CompactionTest.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/CompactionTest.java index beac55d6cc3..ccf6fd43756 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/CompactionTest.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/CompactionTest.java @@ -51,6 +51,7 @@ import java.util.Set; import java.util.concurrent.CompletableFuture; import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import org.apache.bookkeeper.bookie.LedgerDirsManager.NoWritableLedgerDirException; @@ -573,6 +574,11 @@ public void testMinorCompactionWithEntryLogPerLedgerEnabled() throws Exception { bkc.deleteLedger(lhs[1].getId()); bkc.deleteLedger(lhs[2].getId()); + // Need to wait until entry log 3 gets flushed before initiating GC to satisfy assertions. + while (!getGCThread().entryLogger.isFlushedEntryLog(3L)) { + TimeUnit.MILLISECONDS.sleep(100); + } + LOG.info("Finished deleting the ledgers contains most entries."); getGCThread().triggerGC(true, false, false).get(); @@ -589,6 +595,23 @@ public void testMinorCompactionWithEntryLogPerLedgerEnabled() throws Exception { + tmpDirs.get(0), TestUtils.hasAllLogFiles(tmpDirs.get(0), 0, 1, 4)); assertTrue("Found entry log files [2, 3].log that should have been compacted in ledgerDirectory: " + tmpDirs.get(0), TestUtils.hasNoneLogFiles(tmpDirs.get(0), 2, 3)); + + // Now, let's mark E1 as flushed, as its ledger L1 has been deleted already. In this case, the GC algorithm + // should consider it for deletion. + getGCThread().entryLogger.recentlyCreatedEntryLogsStatus.flushRotatedEntryLog(1L); + getGCThread().triggerGC(true, false, false).get(); + assertTrue("Found entry log file 1.log that should have been compacted in ledgerDirectory: " + + tmpDirs.get(0), TestUtils.hasNoneLogFiles(tmpDirs.get(0), 1)); + + // Once removed the ledger L0, then deleting E0 is fine (only if it has been flushed). + bkc.deleteLedger(lhs[0].getId()); + getGCThread().triggerGC(true, false, false).get(); + assertTrue("Found entry log file 0.log that should not have been compacted in ledgerDirectory: " + + tmpDirs.get(0), TestUtils.hasAllLogFiles(tmpDirs.get(0), 0)); + getGCThread().entryLogger.recentlyCreatedEntryLogsStatus.flushRotatedEntryLog(0L); + getGCThread().triggerGC(true, false, false).get(); + assertTrue("Found entry log file 0.log that should have been compacted in ledgerDirectory: " + + tmpDirs.get(0), TestUtils.hasNoneLogFiles(tmpDirs.get(0), 0)); } @Test