From ea2b9b45912070f4160f29dccc7c9325ba358f16 Mon Sep 17 00:00:00 2001 From: rdhabalia Date: Wed, 6 Feb 2019 13:45:59 -0800 Subject: [PATCH] [bk-gc] Prevent OOM: Restrict number of entry-loggers extraction --- .../bookie/GarbageCollectorThread.java | 133 +++++++++++------- .../bookkeeper/conf/ServerConfiguration.java | 20 +++ .../bookkeeper/bookie/CompactionTest.java | 4 +- 3 files changed, 106 insertions(+), 51 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 006592646c0..cf5ef4f0b52 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 @@ -63,6 +63,13 @@ public class GarbageCollectorThread extends SafeRunnable { // This is how often we want to run the Garbage Collector Thread (in milliseconds). final long gcWaitTime; + // Max number of entry-logger files should be extracted concurrently while + // performing GC + final int maxEntryLoggersPerScan; + // flag to iterate over chunk of of entry-logger files + boolean moreEntryLoggers = true; + long lastIterationLogId = 0; + private final boolean verifyMetadataOnGc; // Compaction parameters boolean enableMinorCompaction = false; @@ -146,6 +153,9 @@ public GarbageCollectorThread(ServerConfiguration conf, this.entryLogger = ledgerStorage.getEntryLogger(); this.ledgerStorage = ledgerStorage; this.gcWaitTime = conf.getGcWaitTime(); + + this.maxEntryLoggersPerScan = conf.getMaxEntryLoggersScanOnGc(); + this.verifyMetadataOnGc = conf.getVerifyMetadataOnGC(); this.numActiveEntryLogs = 0; this.totalEntryLogSize = 0L; @@ -322,55 +332,64 @@ public void runWithFlags(boolean force, boolean suspendMajor, boolean suspendMin } // Recover and clean up previous state if using transactional compaction compactor.cleanUpAndRecover(); + + long startTime = System.currentTimeMillis(); + boolean isIteration = false; + do { + // Extract all of the ledger ID's that comprise all of the entry + // logs + // (except for the current new one which is still being written to). + entryLogMetaMap = extractMetaFromEntryLogs(entryLogMetaMap, isIteration); + + // gc inactive/deleted ledgers + doGcLedgers(); + + // gc entry logs + doGcEntryLogs(); + + if (suspendMajor) { + LOG.info("Disk almost full, suspend major compaction to slow down filling disk."); + } + if (suspendMinor) { + LOG.info("Disk full, suspend minor compaction to slow down filling disk."); + } - // Extract all of the ledger ID's that comprise all of the entry logs - // (except for the current new one which is still being written to). - entryLogMetaMap = extractMetaFromEntryLogs(entryLogMetaMap); - - // gc inactive/deleted ledgers - doGcLedgers(); - - // gc entry logs - doGcEntryLogs(); - - if (suspendMajor) { - LOG.info("Disk almost full, suspend major compaction to slow down filling disk."); - } - if (suspendMinor) { - LOG.info("Disk full, suspend minor compaction to slow down filling disk."); - } - - long curTime = System.currentTimeMillis(); - if (enableMajorCompaction && (!suspendMajor) - && (force || curTime - lastMajorCompactionTime > majorCompactionInterval)) { - // enter major compaction - LOG.info("Enter major compaction, suspendMajor {}", suspendMajor); - majorCompacting.set(true); - doCompactEntryLogs(majorCompactionThreshold); - lastMajorCompactionTime = System.currentTimeMillis(); - // and also move minor compaction time - lastMinorCompactionTime = lastMajorCompactionTime; - gcStats.getMajorCompactionCounter().inc(); - majorCompacting.set(false); - } else if (enableMinorCompaction && (!suspendMinor) - && (force || curTime - lastMinorCompactionTime > minorCompactionInterval)) { - // enter minor compaction - LOG.info("Enter minor compaction, suspendMinor {}", suspendMinor); - minorCompacting.set(true); - doCompactEntryLogs(minorCompactionThreshold); - lastMinorCompactionTime = System.currentTimeMillis(); - gcStats.getMinorCompactionCounter().inc(); - minorCompacting.set(false); - } + long curTime = System.currentTimeMillis(); + if (enableMajorCompaction && (!suspendMajor) + && (force || curTime - lastMajorCompactionTime > majorCompactionInterval)) { + // enter major compaction + LOG.info("Enter major compaction, suspendMajor {}", suspendMajor); + majorCompacting.set(true); + doCompactEntryLogs(majorCompactionThreshold); + lastMajorCompactionTime = System.currentTimeMillis(); + // and also move minor compaction time + lastMinorCompactionTime = lastMajorCompactionTime; + gcStats.getMajorCompactionCounter().inc(); + majorCompacting.set(false); + } else if (enableMinorCompaction && (!suspendMinor) + && (force || curTime - lastMinorCompactionTime > minorCompactionInterval)) { + // enter minor compaction + LOG.info("Enter minor compaction, suspendMinor {}", suspendMinor); + minorCompacting.set(true); + doCompactEntryLogs(minorCompactionThreshold); + lastMinorCompactionTime = System.currentTimeMillis(); + gcStats.getMinorCompactionCounter().inc(); + minorCompacting.set(false); + } - if (force) { - if (forceGarbageCollection.compareAndSet(true, false)) { - LOG.info("{} Set forceGarbageCollection to false after force GC to make it forceGC-able again.", Thread - .currentThread().getName()); + if (force) { + if (forceGarbageCollection.compareAndSet(true, false)) { + LOG.info("{} Set forceGarbageCollection to false after force GC to make it forceGC-able again.", + Thread.currentThread().getName()); + } } - } - gcStats.getGcThreadRuntime().registerSuccessfulEvent( - MathUtils.nowInNano() - threadStart, TimeUnit.NANOSECONDS); + gcStats.getGcThreadRuntime().registerSuccessfulEvent(MathUtils.nowInNano() - threadStart, + TimeUnit.NANOSECONDS); + isIteration = true; + } while (moreEntryLoggers); + + long endTime = System.currentTimeMillis(); + LOG.info("Garbage collector completed in {}", TimeUnit.MILLISECONDS.toSeconds(endTime - startTime)); } /** @@ -526,21 +545,35 @@ protected void compactEntryLog(EntryLogMetadata entryLogMeta) { * Existing EntryLogs to Meta * @throws IOException */ - protected Map extractMetaFromEntryLogs(Map entryLogMetaMap) { + protected Map extractMetaFromEntryLogs(Map entryLogMetaMap, boolean isIteration) { + moreEntryLoggers = false; // 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. long curLogId = entryLogger.getLeastUnflushedLogId(); boolean hasExceptionWhenScan = false; - for (long entryLogId = scannedLogId; entryLogId < curLogId; entryLogId++) { - // Comb the current entry log file if it has not already been extracted. + int entryLogFileCount = 0; + long entryLogId = isIteration ? lastIterationLogId : scannedLogId; + for (; entryLogId < curLogId; entryLogId++, entryLogFileCount++) { + if (entryLogFileCount > maxEntryLoggersPerScan && verifyMetadataOnGc) { + lastIterationLogId = entryLogId; + moreEntryLoggers = true; + LOG.debug("extraction max-entry-logger {}, next iteration starts from {}", entryLogFileCount, + maxEntryLoggersPerScan, entryLogId); + break; + } + // Comb the current entry log file if it has not already been + // extracted. if (entryLogMetaMap.containsKey(entryLogId)) { + entryLogFileCount--; continue; } // check whether log file exists or not - // if it doesn't exist, this log file might have been garbage collected. + // if it doesn't exist, this log file might have been garbage + // collected. if (!entryLogger.logExists(entryLogId)) { + entryLogFileCount--; continue; } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/conf/ServerConfiguration.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/conf/ServerConfiguration.java index f4972f0a781..cb87e2af32c 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/conf/ServerConfiguration.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/conf/ServerConfiguration.java @@ -104,6 +104,7 @@ public class ServerConfiguration extends AbstractConfiguration