From 20e6ecc0eb91adf333b2c25f822632d7ce390379 Mon Sep 17 00:00:00 2001 From: gaozhangmin Date: Mon, 25 May 2026 18:21:50 +0800 Subject: [PATCH 1/5] Reuse GC metadata driver for over-replicated cleanup --- .../bookie/GarbageCollectorThread.java | 5 ++ .../ScanAndCompareGarbageCollector.java | 46 ++++++++++++++--- .../bookie/GarbageCollectorThreadTest.java | 50 +++++++++++++++++++ 3 files changed, 93 insertions(+), 8 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 104e6ccc0ef..09613678708 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 @@ -819,6 +819,11 @@ public synchronized void shutdown() throws InterruptedException { } catch (Exception e) { log.warn().exception(e).log("Failed to close entryLog metadata-map"); } + try { + garbageCollector.close(); + } catch (Exception e) { + log.warn().exception(e).log("Failed to close garbage collector metadata resources"); + } } /** diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/ScanAndCompareGarbageCollector.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/ScanAndCompareGarbageCollector.java index 441627737a5..eb00168f8fd 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/ScanAndCompareGarbageCollector.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/ScanAndCompareGarbageCollector.java @@ -70,7 +70,7 @@ *

*/ @CustomLog -public class ScanAndCompareGarbageCollector implements GarbageCollector { +public class ScanAndCompareGarbageCollector implements GarbageCollector, AutoCloseable { private final LedgerManager ledgerManager; private final CompactableLedgerStorage ledgerStorage; @@ -84,6 +84,8 @@ public class ScanAndCompareGarbageCollector implements GarbageCollector { private StatsLogger statsLogger; private final int maxConcurrentRequests; private final RateLimiter gcMetadataOpRateLimiter; + private MetadataBookieDriver metadataDriver; + private LedgerManagerFactory metadataLedgerManagerFactory; public ScanAndCompareGarbageCollector(LedgerManager ledgerManager, CompactableLedgerStorage ledgerStorage, ServerConfiguration conf, StatsLogger statsLogger) throws IOException { @@ -235,13 +237,7 @@ private Set removeOverReplicatedledgers(Set bkActiveledgers, final G final Set overReplicatedLedgers = Sets.newHashSet(); final Semaphore semaphore = new Semaphore(this.maxConcurrentRequests); final CountDownLatch latch = new CountDownLatch(bkActiveledgers.size()); - // instantiate zookeeper client to initialize ledger manager - - @Cleanup - MetadataBookieDriver metadataDriver = instantiateMetadataDriver(conf, statsLogger); - - @Cleanup - LedgerManagerFactory lmf = metadataDriver.getLedgerManagerFactory(); + LedgerManagerFactory lmf = getOrCreateMetadataLedgerManagerFactory(); @Cleanup LedgerUnderreplicationManager lum = lmf.newLedgerUnderreplicationManager(); @@ -324,6 +320,22 @@ private Set removeOverReplicatedledgers(Set bkActiveledgers, final G return overReplicatedLedgers; } + private synchronized LedgerManagerFactory getOrCreateMetadataLedgerManagerFactory() throws BookieException { + if (metadataLedgerManagerFactory != null) { + return metadataLedgerManagerFactory; + } + + metadataDriver = instantiateMetadataDriver(conf, statsLogger); + try { + metadataLedgerManagerFactory = metadataDriver.getLedgerManagerFactory(); + } catch (MetadataException me) { + closeMetadataDriver(); + throw new BookieException.MetadataStoreException( + "Failed to initialize ledger manager factory for over-replicated ledger GC", me); + } + return metadataLedgerManagerFactory; + } + private static MetadataBookieDriver instantiateMetadataDriver(ServerConfiguration conf, StatsLogger statsLogger) throws BookieException { try { @@ -340,6 +352,24 @@ private static MetadataBookieDriver instantiateMetadataDriver(ServerConfiguratio } } + @Override + public synchronized void close() { + closeMetadataDriver(); + } + + private void closeMetadataDriver() { + metadataLedgerManagerFactory = null; + if (metadataDriver != null) { + try { + metadataDriver.close(); + } catch (Exception e) { + log.warn().exception(e).log("Failed to close metadata driver used for over-replicated ledger GC"); + } finally { + metadataDriver = null; + } + } + } + private boolean isNotBookieIncludedInLedgerEnsembles(Versioned metadata) { // do not delete a ledger that is not closed, since the ensemble might // change again and include the current bookie while we are deleting it diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/GarbageCollectorThreadTest.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/GarbageCollectorThreadTest.java index ad2f0c95480..0821248f490 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/GarbageCollectorThreadTest.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/GarbageCollectorThreadTest.java @@ -35,19 +35,28 @@ import static org.junit.Assert.assertTrue; import static org.mockito.Mockito.any; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.mockStatic; import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import static org.mockito.MockitoAnnotations.openMocks; import io.github.merlimat.slog.Event; import io.github.merlimat.slog.Logger; import java.io.File; +import java.lang.reflect.Field; +import java.lang.reflect.Method; +import java.net.URI; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.atomic.AtomicBoolean; import org.apache.bookkeeper.bookie.storage.EntryLogger; import org.apache.bookkeeper.conf.ServerConfiguration; import org.apache.bookkeeper.conf.TestBKConfiguration; import org.apache.bookkeeper.meta.LedgerManager; +import org.apache.bookkeeper.meta.LedgerManagerFactory; +import org.apache.bookkeeper.meta.MetadataBookieDriver; +import org.apache.bookkeeper.meta.MetadataDrivers; import org.apache.bookkeeper.meta.MockLedgerManager; import org.apache.bookkeeper.stats.NullStatsLogger; import org.apache.bookkeeper.stats.StatsLogger; @@ -58,6 +67,7 @@ import org.junit.Test; import org.mockito.InjectMocks; import org.mockito.Mock; +import org.mockito.MockedStatic; import org.mockito.Spy; /** @@ -145,6 +155,46 @@ public void testCalculateUsageBucket() { Assert.assertEquals("Incorrect number of items", items + 1, sum); } + @Test + public void testOverreplicatedLedgerGcReusesMetadataDriverUntilClosed() throws Exception { + ServerConfiguration bkConf = TestBKConfiguration.newServerConfiguration() + .setAllowLoopback(true) + .setMetadataServiceUri("zk://127.0.0.1/ledgers"); + MetadataBookieDriver metadataDriver = mock(MetadataBookieDriver.class); + LedgerManagerFactory lmf = mock(LedgerManagerFactory.class); + when(metadataDriver.getLedgerManagerFactory()).thenReturn(lmf); + + try (MockedStatic metadataDrivers = mockStatic(MetadataDrivers.class)) { + metadataDrivers.when(() -> MetadataDrivers.getBookieDriver(any(URI.class))).thenReturn(metadataDriver); + ScanAndCompareGarbageCollector garbageCollector = new ScanAndCompareGarbageCollector( + ledgerManager, ledgerStorage, bkConf, NullStatsLogger.INSTANCE); + + Assert.assertSame(lmf, getOrCreateMetadataLedgerManagerFactory(garbageCollector)); + Assert.assertSame(lmf, getOrCreateMetadataLedgerManagerFactory(garbageCollector)); + metadataDrivers.verify(() -> MetadataDrivers.getBookieDriver(any(URI.class)), times(1)); + + garbageCollector.close(); + + verify(metadataDriver).close(); + Assert.assertNull(getMetadataDriver(garbageCollector)); + } + } + + private static LedgerManagerFactory getOrCreateMetadataLedgerManagerFactory( + ScanAndCompareGarbageCollector garbageCollector) throws Exception { + Method method = ScanAndCompareGarbageCollector.class.getDeclaredMethod( + "getOrCreateMetadataLedgerManagerFactory"); + method.setAccessible(true); + return (LedgerManagerFactory) method.invoke(garbageCollector); + } + + private static MetadataBookieDriver getMetadataDriver(ScanAndCompareGarbageCollector garbageCollector) + throws Exception { + Field field = ScanAndCompareGarbageCollector.class.getDeclaredField("metadataDriver"); + field.setAccessible(true); + return (MetadataBookieDriver) field.get(garbageCollector); + } + @Test public void testExtractMetaFromEntryLogsLegacy() throws Exception { File ledgerDir = tmpDirs.createNew("testExtractMeta", "ledgers"); From a686ca0e78548750ff3fda0f85de12a1f9596d73 Mon Sep 17 00:00:00 2001 From: gaozhangmin Date: Wed, 27 May 2026 15:27:19 +0800 Subject: [PATCH 2/5] fix exception --- .../bookie/ScanAndCompareGarbageCollector.java | 10 ++-------- 1 file changed, 2 insertions(+), 8 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/ScanAndCompareGarbageCollector.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/ScanAndCompareGarbageCollector.java index eb00168f8fd..8c7ae7a55c2 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/ScanAndCompareGarbageCollector.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/ScanAndCompareGarbageCollector.java @@ -320,19 +320,13 @@ private Set removeOverReplicatedledgers(Set bkActiveledgers, final G return overReplicatedLedgers; } - private synchronized LedgerManagerFactory getOrCreateMetadataLedgerManagerFactory() throws BookieException { + private synchronized LedgerManagerFactory getOrCreateMetadataLedgerManagerFactory() throws Exception { if (metadataLedgerManagerFactory != null) { return metadataLedgerManagerFactory; } metadataDriver = instantiateMetadataDriver(conf, statsLogger); - try { - metadataLedgerManagerFactory = metadataDriver.getLedgerManagerFactory(); - } catch (MetadataException me) { - closeMetadataDriver(); - throw new BookieException.MetadataStoreException( - "Failed to initialize ledger manager factory for over-replicated ledger GC", me); - } + metadataLedgerManagerFactory = metadataDriver.getLedgerManagerFactory(); return metadataLedgerManagerFactory; } From 5b91bfb770d19fac7853ae89e44c93d2f4ac1105 Mon Sep 17 00:00:00 2001 From: gaozhangmin Date: Wed, 27 May 2026 15:32:23 +0800 Subject: [PATCH 3/5] fix error --- .../bookkeeper/bookie/GarbageCollectorThread.java | 2 +- .../bookie/ScanAndCompareGarbageCollector.java | 12 ++++-------- .../bookie/GarbageCollectorThreadTest.java | 2 +- 3 files changed, 6 insertions(+), 10 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 09613678708..a1adaf92cb0 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 @@ -820,7 +820,7 @@ public synchronized void shutdown() throws InterruptedException { log.warn().exception(e).log("Failed to close entryLog metadata-map"); } try { - garbageCollector.close(); + garbageCollector.closeMetadataDriver(); } catch (Exception e) { log.warn().exception(e).log("Failed to close garbage collector metadata resources"); } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/ScanAndCompareGarbageCollector.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/ScanAndCompareGarbageCollector.java index 8c7ae7a55c2..552ae7b4f1e 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/ScanAndCompareGarbageCollector.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/ScanAndCompareGarbageCollector.java @@ -70,7 +70,7 @@ *

*/ @CustomLog -public class ScanAndCompareGarbageCollector implements GarbageCollector, AutoCloseable { +public class ScanAndCompareGarbageCollector implements GarbageCollector { private final LedgerManager ledgerManager; private final CompactableLedgerStorage ledgerStorage; @@ -346,19 +346,15 @@ private static MetadataBookieDriver instantiateMetadataDriver(ServerConfiguratio } } - @Override - public synchronized void close() { - closeMetadataDriver(); - } - - private void closeMetadataDriver() { - metadataLedgerManagerFactory = null; + public void closeMetadataDriver() { if (metadataDriver != null) { try { + metadataLedgerManagerFactory.close(); metadataDriver.close(); } catch (Exception e) { log.warn().exception(e).log("Failed to close metadata driver used for over-replicated ledger GC"); } finally { + metadataLedgerManagerFactory = null; metadataDriver = null; } } diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/GarbageCollectorThreadTest.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/GarbageCollectorThreadTest.java index 0821248f490..0809d34abfb 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/GarbageCollectorThreadTest.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/GarbageCollectorThreadTest.java @@ -173,7 +173,7 @@ public void testOverreplicatedLedgerGcReusesMetadataDriverUntilClosed() throws E Assert.assertSame(lmf, getOrCreateMetadataLedgerManagerFactory(garbageCollector)); metadataDrivers.verify(() -> MetadataDrivers.getBookieDriver(any(URI.class)), times(1)); - garbageCollector.close(); + garbageCollector.closeMetadataDriver(); verify(metadataDriver).close(); Assert.assertNull(getMetadataDriver(garbageCollector)); From ce2225e52fbb1758b37d481047601fbc9aa91e99 Mon Sep 17 00:00:00 2001 From: gaozhangmin Date: Mon, 13 Jul 2026 16:51:35 +0800 Subject: [PATCH 4/5] Reuse server ledger manager factory for GC --- .../apache/bookkeeper/bookie/BookieImpl.java | 1 + .../bookkeeper/bookie/BookieResources.java | 3 + .../bookie/GarbageCollectorThread.java | 33 ++++++++-- .../bookie/InterleavedLedgerStorage.java | 9 ++- .../bookkeeper/bookie/LedgerStorage.java | 2 + .../ScanAndCompareGarbageCollector.java | 64 ++++--------------- .../bookie/SortedLedgerStorage.java | 6 ++ .../bookie/storage/ldb/DbLedgerStorage.java | 20 ++++-- .../ldb/SingleDirectoryDbLedgerStorage.java | 6 +- .../bookkeeper/server/EmbeddedServer.java | 4 +- .../bookkeeper/util/LocalBookKeeper.java | 2 +- .../bookie/GarbageCollectorThreadTest.java | 50 +++------------ .../bookie/GcOverreplicatedLedgerTest.java | 31 ++------- .../bookkeeper/bookie/TestBookieImpl.java | 2 +- .../ldb/DbLedgerStorageWriteCacheTest.java | 17 ++--- .../test/BookKeeperClusterTestCase.java | 2 +- 16 files changed, 103 insertions(+), 149 deletions(-) diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/BookieImpl.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/BookieImpl.java index 362a20129a1..cffd1b23de1 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/BookieImpl.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/BookieImpl.java @@ -363,6 +363,7 @@ public static LedgerStorage mountLedgerStorageOffline(ServerConfiguration conf, if (null == ledgerStorage) { ledgerStorage = BookieResources.createLedgerStorage(conf, null, + null, ledgerDirsManager, indexDirsManager, statsLogger, diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/BookieResources.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/BookieResources.java index 919f95af6af..fb95a8017b3 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/BookieResources.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/BookieResources.java @@ -28,6 +28,7 @@ import org.apache.bookkeeper.common.allocator.ByteBufAllocatorWithOomHandler; import org.apache.bookkeeper.conf.ServerConfiguration; import org.apache.bookkeeper.meta.LedgerManager; +import org.apache.bookkeeper.meta.LedgerManagerFactory; import org.apache.bookkeeper.meta.MetadataBookieDriver; import org.apache.bookkeeper.meta.MetadataDrivers; import org.apache.bookkeeper.meta.exceptions.MetadataException; @@ -95,6 +96,7 @@ public static LedgerDirsManager createIndexDirsManager(ServerConfiguration conf, public static LedgerStorage createLedgerStorage(ServerConfiguration conf, LedgerManager ledgerManager, + LedgerManagerFactory ledgerManagerFactory, LedgerDirsManager ledgerDirsManager, LedgerDirsManager indexDirsManager, StatsLogger statsLogger, @@ -104,6 +106,7 @@ public static LedgerStorage createLedgerStorage(ServerConfiguration conf, log.info().attr("ledgerStorageClass", ledgerStorageClass).log("Using ledger storage"); LedgerStorage storage = LedgerStorageFactory.createLedgerStorage(ledgerStorageClass); + storage.setLedgerManagerFactory(ledgerManagerFactory); storage.initialize(conf, ledgerManager, ledgerDirsManager, indexDirsManager, statsLogger, allocator); storage.setCheckpointSource(CheckpointSource.DEFAULT); return storage; 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 a1adaf92cb0..81bac734b19 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 @@ -49,6 +49,7 @@ import org.apache.bookkeeper.common.util.MathUtils; import org.apache.bookkeeper.conf.ServerConfiguration; import org.apache.bookkeeper.meta.LedgerManager; +import org.apache.bookkeeper.meta.LedgerManagerFactory; import org.apache.bookkeeper.stats.StatsLogger; import org.apache.commons.lang3.mutable.MutableBoolean; import org.apache.commons.lang3.mutable.MutableLong; @@ -147,7 +148,17 @@ public GarbageCollectorThread(ServerConfiguration conf, LedgerManager ledgerMana final CompactableLedgerStorage ledgerStorage, EntryLogger entryLogger, StatsLogger statsLogger) throws IOException { - this(conf, ledgerManager, ledgerDirsManager, ledgerStorage, entryLogger, statsLogger, newExecutor()); + this(conf, ledgerManager, null, ledgerDirsManager, ledgerStorage, entryLogger, statsLogger, newExecutor()); + } + + public GarbageCollectorThread(ServerConfiguration conf, LedgerManager ledgerManager, + LedgerManagerFactory ledgerManagerFactory, + final LedgerDirsManager ledgerDirsManager, + final CompactableLedgerStorage ledgerStorage, + EntryLogger entryLogger, + StatsLogger statsLogger) throws IOException { + this(conf, ledgerManager, ledgerManagerFactory, ledgerDirsManager, ledgerStorage, entryLogger, statsLogger, + newExecutor()); } @VisibleForTesting @@ -170,6 +181,18 @@ public GarbageCollectorThread(ServerConfiguration conf, StatsLogger statsLogger, ScheduledExecutorService gcExecutor) throws IOException { + this(conf, ledgerManager, null, ledgerDirsManager, ledgerStorage, entryLogger, statsLogger, gcExecutor); + } + + public GarbageCollectorThread(ServerConfiguration conf, + LedgerManager ledgerManager, + LedgerManagerFactory ledgerManagerFactory, + final LedgerDirsManager ledgerDirsManager, + final CompactableLedgerStorage ledgerStorage, + EntryLogger entryLogger, + StatsLogger statsLogger, + ScheduledExecutorService gcExecutor) + throws IOException { this.gcExecutor = gcExecutor; this.conf = conf; @@ -184,7 +207,8 @@ public GarbageCollectorThread(ServerConfiguration conf, this.totalEntryLogSize = 0L; this.entryLogCompactRatio = 0.0; this.currentEntryLogUsageBuckets = new int[ENTRY_LOG_USAGE_SEGMENT_COUNT]; - this.garbageCollector = new ScanAndCompareGarbageCollector(ledgerManager, ledgerStorage, conf, statsLogger); + this.garbageCollector = new ScanAndCompareGarbageCollector(ledgerManager, ledgerStorage, + ledgerManagerFactory, conf, statsLogger); this.gcStats = new GarbageCollectorStats( statsLogger, () -> numActiveEntryLogs, @@ -819,11 +843,6 @@ public synchronized void shutdown() throws InterruptedException { } catch (Exception e) { log.warn().exception(e).log("Failed to close entryLog metadata-map"); } - try { - garbageCollector.closeMetadataDriver(); - } catch (Exception e) { - log.warn().exception(e).log("Failed to close garbage collector metadata resources"); - } } /** diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/InterleavedLedgerStorage.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/InterleavedLedgerStorage.java index fa1b724622d..74b464c8f7c 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/InterleavedLedgerStorage.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/InterleavedLedgerStorage.java @@ -62,6 +62,7 @@ import org.apache.bookkeeper.common.util.Watcher; import org.apache.bookkeeper.conf.ServerConfiguration; import org.apache.bookkeeper.meta.LedgerManager; +import org.apache.bookkeeper.meta.LedgerManagerFactory; import org.apache.bookkeeper.proto.BookieProtocol; import org.apache.bookkeeper.stats.Counter; import org.apache.bookkeeper.stats.OpStatsLogger; @@ -100,6 +101,7 @@ public class InterleavedLedgerStorage implements CompactableLedgerStorage, Entry // contain any active ledgers in them; and compacts the entry logs that // has lower remaining percentage to reclaim disk space. GarbageCollectorThread gcThread; + private LedgerManagerFactory ledgerManagerFactory; // this indicates that a write has happened since the last flush private final AtomicBoolean somethingWritten = new AtomicBoolean(false); @@ -163,6 +165,11 @@ void initializeWithEntryLogListener(ServerConfiguration conf, @Override public void setStateManager(StateManager stateManager) {} + @Override + public void setLedgerManagerFactory(LedgerManagerFactory ledgerManagerFactory) { + this.ledgerManagerFactory = ledgerManagerFactory; + } + @Override public void setCheckpointSource(CheckpointSource checkpointSource) { this.checkpointSource = checkpointSource; @@ -185,7 +192,7 @@ public void initializeWithEntryLogger(ServerConfiguration conf, this.entryLogger.addListener(this); ledgerCache = new LedgerCacheImpl(conf, activeLedgers, null == indexDirsManager ? ledgerDirsManager : indexDirsManager, statsLogger); - gcThread = new GarbageCollectorThread(conf, ledgerManager, ledgerDirsManager, + gcThread = new GarbageCollectorThread(conf, ledgerManager, ledgerManagerFactory, ledgerDirsManager, this, entryLogger, statsLogger.scope("gc")); ledgerDirsManager.addLedgerDirsListener(getLedgerDirsListener()); // Expose Stats diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/LedgerStorage.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/LedgerStorage.java index 6eca6e00108..23ef0d41013 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/LedgerStorage.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/LedgerStorage.java @@ -36,6 +36,7 @@ import org.apache.bookkeeper.common.util.Watcher; import org.apache.bookkeeper.conf.ServerConfiguration; import org.apache.bookkeeper.meta.LedgerManager; +import org.apache.bookkeeper.meta.LedgerManagerFactory; import org.apache.bookkeeper.stats.StatsLogger; /** @@ -59,6 +60,7 @@ void initialize(ServerConfiguration conf, throws IOException; void setStateManager(StateManager stateManager); + default void setLedgerManagerFactory(LedgerManagerFactory ledgerManagerFactory) {} void setCheckpointSource(CheckpointSource checkpointSource); void setCheckpointer(Checkpointer checkpointer); diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/ScanAndCompareGarbageCollector.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/ScanAndCompareGarbageCollector.java index 552ae7b4f1e..3aface9def0 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/ScanAndCompareGarbageCollector.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/ScanAndCompareGarbageCollector.java @@ -26,7 +26,6 @@ import com.google.common.collect.Sets; import com.google.common.util.concurrent.RateLimiter; import java.io.IOException; -import java.net.URI; import java.util.List; import java.util.NavigableSet; import java.util.Set; @@ -47,13 +46,9 @@ import org.apache.bookkeeper.meta.LedgerManager.LedgerRangeIterator; import org.apache.bookkeeper.meta.LedgerManagerFactory; import org.apache.bookkeeper.meta.LedgerUnderreplicationManager; -import org.apache.bookkeeper.meta.MetadataBookieDriver; -import org.apache.bookkeeper.meta.MetadataDrivers; -import org.apache.bookkeeper.meta.exceptions.MetadataException; import org.apache.bookkeeper.net.BookieId; import org.apache.bookkeeper.stats.StatsLogger; import org.apache.bookkeeper.versioning.Versioned; -import org.apache.commons.configuration2.ex.ConfigurationException; /** * Garbage collector implementation using scan and compare. @@ -81,18 +76,22 @@ public class ScanAndCompareGarbageCollector implements GarbageCollector { private long lastOverReplicatedLedgerGcTimeMillis; private final boolean verifyMetadataOnGc; private int activeLedgerCounter; - private StatsLogger statsLogger; private final int maxConcurrentRequests; private final RateLimiter gcMetadataOpRateLimiter; - private MetadataBookieDriver metadataDriver; - private LedgerManagerFactory metadataLedgerManagerFactory; + private final LedgerManagerFactory ledgerManagerFactory; public ScanAndCompareGarbageCollector(LedgerManager ledgerManager, CompactableLedgerStorage ledgerStorage, ServerConfiguration conf, StatsLogger statsLogger) throws IOException { + this(ledgerManager, ledgerStorage, null, conf, statsLogger); + } + + public ScanAndCompareGarbageCollector(LedgerManager ledgerManager, CompactableLedgerStorage ledgerStorage, + LedgerManagerFactory ledgerManagerFactory, ServerConfiguration conf, StatsLogger statsLogger) + throws IOException { this.ledgerManager = ledgerManager; this.ledgerStorage = ledgerStorage; + this.ledgerManagerFactory = ledgerManagerFactory; this.conf = conf; - this.statsLogger = statsLogger; this.selfBookieAddress = BookieImpl.getBookieId(conf); this.gcOverReplicatedLedgerIntervalMillis = conf.getGcOverreplicatedLedgerWaitTimeMillis(); @@ -237,10 +236,13 @@ private Set removeOverReplicatedledgers(Set bkActiveledgers, final G final Set overReplicatedLedgers = Sets.newHashSet(); final Semaphore semaphore = new Semaphore(this.maxConcurrentRequests); final CountDownLatch latch = new CountDownLatch(bkActiveledgers.size()); - LedgerManagerFactory lmf = getOrCreateMetadataLedgerManagerFactory(); + if (ledgerManagerFactory == null) { + log.warn("Skipping over-replicated ledger GC because LedgerManagerFactory is not available."); + return overReplicatedLedgers; + } @Cleanup - LedgerUnderreplicationManager lum = lmf.newLedgerUnderreplicationManager(); + LedgerUnderreplicationManager lum = ledgerManagerFactory.newLedgerUnderreplicationManager(); for (final Long ledgerId : bkActiveledgers) { try { @@ -320,46 +322,6 @@ private Set removeOverReplicatedledgers(Set bkActiveledgers, final G return overReplicatedLedgers; } - private synchronized LedgerManagerFactory getOrCreateMetadataLedgerManagerFactory() throws Exception { - if (metadataLedgerManagerFactory != null) { - return metadataLedgerManagerFactory; - } - - metadataDriver = instantiateMetadataDriver(conf, statsLogger); - metadataLedgerManagerFactory = metadataDriver.getLedgerManagerFactory(); - return metadataLedgerManagerFactory; - } - - private static MetadataBookieDriver instantiateMetadataDriver(ServerConfiguration conf, StatsLogger statsLogger) - throws BookieException { - try { - String metadataServiceUriStr = conf.getMetadataServiceUri(); - MetadataBookieDriver driver = MetadataDrivers.getBookieDriver(URI.create(metadataServiceUriStr)); - driver.initialize( - conf, - statsLogger); - return driver; - } catch (MetadataException me) { - throw new BookieException.MetadataStoreException("Failed to initialize metadata bookie driver", me); - } catch (ConfigurationException e) { - throw new BookieException.BookieIllegalOpException(e); - } - } - - public void closeMetadataDriver() { - if (metadataDriver != null) { - try { - metadataLedgerManagerFactory.close(); - metadataDriver.close(); - } catch (Exception e) { - log.warn().exception(e).log("Failed to close metadata driver used for over-replicated ledger GC"); - } finally { - metadataLedgerManagerFactory = null; - metadataDriver = null; - } - } - } - private boolean isNotBookieIncludedInLedgerEnsembles(Versioned metadata) { // do not delete a ledger that is not closed, since the ensemble might // change again and include the current bookie while we are deleting it diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/SortedLedgerStorage.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/SortedLedgerStorage.java index bc3d586008c..3cc8238d657 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/SortedLedgerStorage.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/SortedLedgerStorage.java @@ -38,6 +38,7 @@ import org.apache.bookkeeper.common.util.Watcher; import org.apache.bookkeeper.conf.ServerConfiguration; import org.apache.bookkeeper.meta.LedgerManager; +import org.apache.bookkeeper.meta.LedgerManagerFactory; import org.apache.bookkeeper.proto.BookieProtocol; import org.apache.bookkeeper.stats.StatsLogger; import org.apache.bookkeeper.util.IteratorUtility; @@ -69,6 +70,11 @@ protected SortedLedgerStorage(InterleavedLedgerStorage ils) { interleavedLedgerStorage = ils; } + @Override + public void setLedgerManagerFactory(LedgerManagerFactory ledgerManagerFactory) { + interleavedLedgerStorage.setLedgerManagerFactory(ledgerManagerFactory); + } + @Override public void initialize(ServerConfiguration conf, LedgerManager ledgerManager, diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/DbLedgerStorage.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/DbLedgerStorage.java index 0f672de085d..e9b21b9f753 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/DbLedgerStorage.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/DbLedgerStorage.java @@ -65,6 +65,7 @@ import org.apache.bookkeeper.common.util.nativeio.NativeIOImpl; import org.apache.bookkeeper.conf.ServerConfiguration; import org.apache.bookkeeper.meta.LedgerManager; +import org.apache.bookkeeper.meta.LedgerManagerFactory; import org.apache.bookkeeper.stats.Gauge; import org.apache.bookkeeper.stats.NullStatsLogger; import org.apache.bookkeeper.stats.StatsLogger; @@ -123,6 +124,7 @@ public class DbLedgerStorage implements LedgerStorage { private static final long STORAGE_FLAGS_KEY = 0L; private int numberOfDirs; private List ledgerStorageList; + private LedgerManagerFactory ledgerManagerFactory; private ExecutorService entryLoggerWriteExecutor = null; private ExecutorService entryLoggerFlushExecutor = null; @@ -238,7 +240,7 @@ public void initialize(ServerConfiguration conf, LedgerManager ledgerManager, Le } else { entrylogger = new DefaultEntryLogger(conf, ldm, null, statsLogger, allocator); } - ledgerStorageList.add(newSingleDirectoryDbLedgerStorage(conf, ledgerManager, ldm, + ledgerStorageList.add(newSingleDirectoryDbLedgerStorage(conf, ledgerManager, ledgerManagerFactory, ldm, idm, entrylogger, statsLogger, perDirectoryWriteCacheSize, perDirectoryReadCacheSize, @@ -279,15 +281,21 @@ public Long getSample() { @VisibleForTesting protected SingleDirectoryDbLedgerStorage newSingleDirectoryDbLedgerStorage(ServerConfiguration conf, - LedgerManager ledgerManager, LedgerDirsManager ledgerDirsManager, LedgerDirsManager indexDirsManager, - EntryLogger entryLogger, StatsLogger statsLogger, long writeCacheSize, long readCacheSize, - int readAheadCacheBatchSize, long readAheadCacheBatchBytesSize) + LedgerManager ledgerManager, LedgerManagerFactory ledgerManagerFactory, LedgerDirsManager ledgerDirsManager, + LedgerDirsManager indexDirsManager, EntryLogger entryLogger, StatsLogger statsLogger, long writeCacheSize, + long readCacheSize, int readAheadCacheBatchSize, long readAheadCacheBatchBytesSize) throws IOException { - return new SingleDirectoryDbLedgerStorage(conf, ledgerManager, ledgerDirsManager, indexDirsManager, entryLogger, - statsLogger, allocator, writeCacheSize, readCacheSize, + return new SingleDirectoryDbLedgerStorage(conf, ledgerManager, ledgerManagerFactory, ledgerDirsManager, + indexDirsManager, entryLogger, statsLogger, allocator, writeCacheSize, + readCacheSize, readAheadCacheBatchSize, readAheadCacheBatchBytesSize); } + @Override + public void setLedgerManagerFactory(LedgerManagerFactory ledgerManagerFactory) { + this.ledgerManagerFactory = ledgerManagerFactory; + } + @Override public void setStateManager(StateManager stateManager) { ledgerStorageList.forEach(s -> s.setStateManager(stateManager)); diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/SingleDirectoryDbLedgerStorage.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/SingleDirectoryDbLedgerStorage.java index 6778de2f824..af65351f039 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/SingleDirectoryDbLedgerStorage.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/bookie/storage/ldb/SingleDirectoryDbLedgerStorage.java @@ -72,6 +72,7 @@ import org.apache.bookkeeper.common.util.Watcher; import org.apache.bookkeeper.conf.ServerConfiguration; import org.apache.bookkeeper.meta.LedgerManager; +import org.apache.bookkeeper.meta.LedgerManagerFactory; import org.apache.bookkeeper.proto.BookieProtocol; import org.apache.bookkeeper.stats.Counter; import org.apache.bookkeeper.stats.OpStatsLogger; @@ -152,6 +153,7 @@ protected Thread newThread(Runnable r, String name) { private final String indexBaseDir; public SingleDirectoryDbLedgerStorage(ServerConfiguration conf, LedgerManager ledgerManager, + LedgerManagerFactory ledgerManagerFactory, LedgerDirsManager ledgerDirsManager, LedgerDirsManager indexDirsManager, EntryLogger entryLogger, StatsLogger statsLogger, ByteBufAllocator allocator, long writeCacheSize, long readCacheSize, int readAheadCacheBatchSize, @@ -210,8 +212,8 @@ public SingleDirectoryDbLedgerStorage(ServerConfiguration conf, LedgerManager le TransientLedgerInfo.LEDGER_INFO_CACHING_TIME_MINUTES, TimeUnit.MINUTES); this.entryLogger = entryLogger; - gcThread = new GarbageCollectorThread(conf, - ledgerManager, ledgerDirsManager, this, entryLogger, ledgerIndexDirStatsLogger); + gcThread = new GarbageCollectorThread(conf, ledgerManager, ledgerManagerFactory, ledgerDirsManager, this, + entryLogger, ledgerIndexDirStatsLogger); dbLedgerStorageStats = new DbLedgerStorageStats( ledgerIndexDirStatsLogger, diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/server/EmbeddedServer.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/server/EmbeddedServer.java index 4a957d2dcc6..17980dba197 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/server/EmbeddedServer.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/server/EmbeddedServer.java @@ -387,7 +387,7 @@ public EmbeddedServer build() throws Exception { new RxSchedulerLifecycleComponent("rx-scheduler", conf, bookieStats, rxScheduler, rxExecutor)); - storage = BookieResources.createLedgerStorage(conf.getServerConf(), ledgerManager, + storage = BookieResources.createLedgerStorage(conf.getServerConf(), ledgerManager, ledgerManagerFactory, ledgerDirsManager, indexDirsManager, bookieStats, allocator); EntryCopier copier = new EntryCopierImpl(bookieId, @@ -413,7 +413,7 @@ public EmbeddedServer build() throws Exception { registrationManager); cookieValidation.checkCookies(storageDirectoriesFromConf(conf.getServerConf())); // storage should be created after legacy validation or it will fail (it would find ledger dirs) - storage = BookieResources.createLedgerStorage(conf.getServerConf(), ledgerManager, + storage = BookieResources.createLedgerStorage(conf.getServerConf(), ledgerManager, ledgerManagerFactory, ledgerDirsManager, indexDirsManager, bookieStats, allocator); } diff --git a/bookkeeper-server/src/main/java/org/apache/bookkeeper/util/LocalBookKeeper.java b/bookkeeper-server/src/main/java/org/apache/bookkeeper/util/LocalBookKeeper.java index d0645d7d005..d5eb94b5f0d 100644 --- a/bookkeeper-server/src/main/java/org/apache/bookkeeper/util/LocalBookKeeper.java +++ b/bookkeeper-server/src/main/java/org/apache/bookkeeper/util/LocalBookKeeper.java @@ -531,7 +531,7 @@ private class LocalBookie { LedgerDirsManager indexDirsManager = BookieResources.createIndexDirsManager( conf, diskChecker, NullStatsLogger.INSTANCE, ledgerDirsManager); LedgerStorage storage = BookieResources.createLedgerStorage( - conf, ledgerManager, ledgerDirsManager, indexDirsManager, + conf, ledgerManager, lmFactory, ledgerDirsManager, indexDirsManager, NullStatsLogger.INSTANCE, allocator); CookieValidation cookieValidation = new LegacyCookieValidation(conf, registrationManager); diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/GarbageCollectorThreadTest.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/GarbageCollectorThreadTest.java index 0809d34abfb..500806a35c3 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/GarbageCollectorThreadTest.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/GarbageCollectorThreadTest.java @@ -35,9 +35,8 @@ import static org.junit.Assert.assertTrue; import static org.mockito.Mockito.any; import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.mockStatic; +import static org.mockito.Mockito.never; import static org.mockito.Mockito.spy; -import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import static org.mockito.MockitoAnnotations.openMocks; @@ -45,9 +44,6 @@ import io.github.merlimat.slog.Event; import io.github.merlimat.slog.Logger; import java.io.File; -import java.lang.reflect.Field; -import java.lang.reflect.Method; -import java.net.URI; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.atomic.AtomicBoolean; import org.apache.bookkeeper.bookie.storage.EntryLogger; @@ -55,8 +51,6 @@ import org.apache.bookkeeper.conf.TestBKConfiguration; import org.apache.bookkeeper.meta.LedgerManager; import org.apache.bookkeeper.meta.LedgerManagerFactory; -import org.apache.bookkeeper.meta.MetadataBookieDriver; -import org.apache.bookkeeper.meta.MetadataDrivers; import org.apache.bookkeeper.meta.MockLedgerManager; import org.apache.bookkeeper.stats.NullStatsLogger; import org.apache.bookkeeper.stats.StatsLogger; @@ -67,7 +61,6 @@ import org.junit.Test; import org.mockito.InjectMocks; import org.mockito.Mock; -import org.mockito.MockedStatic; import org.mockito.Spy; /** @@ -156,43 +149,16 @@ public void testCalculateUsageBucket() { } @Test - public void testOverreplicatedLedgerGcReusesMetadataDriverUntilClosed() throws Exception { - ServerConfiguration bkConf = TestBKConfiguration.newServerConfiguration() - .setAllowLoopback(true) - .setMetadataServiceUri("zk://127.0.0.1/ledgers"); - MetadataBookieDriver metadataDriver = mock(MetadataBookieDriver.class); + public void testGarbageCollectorThreadDoesNotCloseServerOwnedLedgerManagerFactory() throws Exception { LedgerManagerFactory lmf = mock(LedgerManagerFactory.class); - when(metadataDriver.getLedgerManagerFactory()).thenReturn(lmf); - - try (MockedStatic metadataDrivers = mockStatic(MetadataDrivers.class)) { - metadataDrivers.when(() -> MetadataDrivers.getBookieDriver(any(URI.class))).thenReturn(metadataDriver); - ScanAndCompareGarbageCollector garbageCollector = new ScanAndCompareGarbageCollector( - ledgerManager, ledgerStorage, bkConf, NullStatsLogger.INSTANCE); - - Assert.assertSame(lmf, getOrCreateMetadataLedgerManagerFactory(garbageCollector)); - Assert.assertSame(lmf, getOrCreateMetadataLedgerManagerFactory(garbageCollector)); - metadataDrivers.verify(() -> MetadataDrivers.getBookieDriver(any(URI.class)), times(1)); - - garbageCollector.closeMetadataDriver(); - - verify(metadataDriver).close(); - Assert.assertNull(getMetadataDriver(garbageCollector)); - } - } + File ledgerDir = tmpDirs.createNew("testGcFactoryOwnership", "ledgers"); + GarbageCollectorThread gcThread = new GarbageCollectorThread( + TestBKConfiguration.newServerConfiguration(), new MockLedgerManager(), lmf, newDirsManager(ledgerDir), + new MockLedgerStorage(), newLegacyEntryLogger(20000, ledgerDir), NullStatsLogger.INSTANCE); - private static LedgerManagerFactory getOrCreateMetadataLedgerManagerFactory( - ScanAndCompareGarbageCollector garbageCollector) throws Exception { - Method method = ScanAndCompareGarbageCollector.class.getDeclaredMethod( - "getOrCreateMetadataLedgerManagerFactory"); - method.setAccessible(true); - return (LedgerManagerFactory) method.invoke(garbageCollector); - } + gcThread.shutdown(); - private static MetadataBookieDriver getMetadataDriver(ScanAndCompareGarbageCollector garbageCollector) - throws Exception { - Field field = ScanAndCompareGarbageCollector.class.getDeclaredField("metadataDriver"); - field.setAccessible(true); - return (MetadataBookieDriver) field.get(garbageCollector); + verify(lmf, never()).close(); } @Test diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/GcOverreplicatedLedgerTest.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/GcOverreplicatedLedgerTest.java index 7ef9390f83e..5be38f3353c 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/GcOverreplicatedLedgerTest.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/GcOverreplicatedLedgerTest.java @@ -23,7 +23,6 @@ import com.google.common.collect.Lists; import java.io.IOException; -import java.net.URI; import java.util.Arrays; import java.util.Collection; import java.util.List; @@ -39,15 +38,11 @@ import org.apache.bookkeeper.meta.LedgerManagerFactory; import org.apache.bookkeeper.meta.LedgerManagerTestCase; import org.apache.bookkeeper.meta.LedgerUnderreplicationManager; -import org.apache.bookkeeper.meta.MetadataBookieDriver; -import org.apache.bookkeeper.meta.MetadataDrivers; import org.apache.bookkeeper.meta.ZkLedgerUnderreplicationManager; -import org.apache.bookkeeper.meta.exceptions.MetadataException; import org.apache.bookkeeper.meta.zk.ZKMetadataDriverBase; import org.apache.bookkeeper.net.BookieId; import org.apache.bookkeeper.stats.NullStatsLogger; import org.apache.bookkeeper.util.SnapshotMap; -import org.apache.commons.configuration2.ex.ConfigurationException; import org.apache.zookeeper.ZooDefs; import org.junit.Assert; import org.junit.Before; @@ -90,11 +85,7 @@ public void testGcOverreplicatedLedger() throws Exception { ServerConfiguration bkConf = getBkConf(bookieNotInEnsemble); @Cleanup - final MetadataBookieDriver metadataDriver = instantiateMetadataDriver(bkConf); - @Cleanup - final LedgerManagerFactory lmf = metadataDriver.getLedgerManagerFactory(); - @Cleanup - final LedgerUnderreplicationManager lum = lmf.newLedgerUnderreplicationManager(); + final LedgerUnderreplicationManager lum = ledgerManagerFactory.newLedgerUnderreplicationManager(); Assert.assertFalse(lum.isLedgerBeingReplicated(lh.getId())); @@ -104,7 +95,7 @@ public void testGcOverreplicatedLedger() throws Exception { final CompactableLedgerStorage mockLedgerStorage = new MockLedgerStorage(); final GarbageCollector garbageCollector = new ScanAndCompareGarbageCollector(ledgerManager, mockLedgerStorage, - bkConf, NullStatsLogger.INSTANCE); + ledgerManagerFactory, bkConf, NullStatsLogger.INSTANCE); Thread.sleep(bkConf.getGcOverreplicatedLedgerWaitTimeMillis() + 1); garbageCollector.gc(new GarbageCleaner() { @@ -123,20 +114,6 @@ public void clean(long ledgerId) { Assert.assertFalse(activeLedgers.containsKey(lh.getId())); } - private static MetadataBookieDriver instantiateMetadataDriver(ServerConfiguration conf) - throws BookieException { - try { - final String metadataServiceUriStr = conf.getMetadataServiceUri(); - final MetadataBookieDriver driver = MetadataDrivers.getBookieDriver(URI.create(metadataServiceUriStr)); - driver.initialize(conf, NullStatsLogger.INSTANCE); - return driver; - } catch (MetadataException me) { - throw new BookieException.MetadataStoreException("Failed to initialize metadata bookie driver", me); - } catch (ConfigurationException e) { - throw new BookieException.BookieIllegalOpException(e); - } - } - @Test public void testNoGcOfLedger() throws Exception { LedgerHandle lh = bkc.createLedger(2, 2, DigestType.MAC, "".getBytes()); @@ -155,7 +132,7 @@ public void testNoGcOfLedger() throws Exception { final CompactableLedgerStorage mockLedgerStorage = new MockLedgerStorage(); final GarbageCollector garbageCollector = new ScanAndCompareGarbageCollector(ledgerManager, mockLedgerStorage, - bkConf, NullStatsLogger.INSTANCE); + ledgerManagerFactory, bkConf, NullStatsLogger.INSTANCE); Thread.sleep(bkConf.getGcOverreplicatedLedgerWaitTimeMillis() + 1); garbageCollector.gc(new GarbageCleaner() { @@ -193,7 +170,7 @@ public void testNoGcIfLedgerBeingReplicated() throws Exception { final CompactableLedgerStorage mockLedgerStorage = new MockLedgerStorage(); final GarbageCollector garbageCollector = new ScanAndCompareGarbageCollector(ledgerManager, mockLedgerStorage, - bkConf, NullStatsLogger.INSTANCE); + ledgerManagerFactory, bkConf, NullStatsLogger.INSTANCE); Thread.sleep(bkConf.getGcOverreplicatedLedgerWaitTimeMillis() + 1); garbageCollector.gc(new GarbageCleaner() { diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/TestBookieImpl.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/TestBookieImpl.java index cd0e967b61c..15db780f236 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/TestBookieImpl.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/TestBookieImpl.java @@ -196,7 +196,7 @@ public Resources build(StatsLogger statsLogger) throws Exception { conf, diskChecker, statsLogger, ledgerDirsManager); LedgerStorage storage = BookieResources.createLedgerStorage( - conf, ledgerManager, ledgerDirsManager, indexDirsManager, + conf, ledgerManager, ledgerManagerFactory, ledgerDirsManager, indexDirsManager, statsLogger, UnpooledByteBufAllocator.DEFAULT); return new Resources(conf, diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/storage/ldb/DbLedgerStorageWriteCacheTest.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/storage/ldb/DbLedgerStorageWriteCacheTest.java index 102f7f5addc..4c2158351d9 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/storage/ldb/DbLedgerStorageWriteCacheTest.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/bookie/storage/ldb/DbLedgerStorageWriteCacheTest.java @@ -37,6 +37,7 @@ import org.apache.bookkeeper.conf.ServerConfiguration; import org.apache.bookkeeper.conf.TestBKConfiguration; import org.apache.bookkeeper.meta.LedgerManager; +import org.apache.bookkeeper.meta.LedgerManagerFactory; import org.apache.bookkeeper.stats.StatsLogger; import org.junit.After; import org.junit.Before; @@ -54,23 +55,23 @@ private static class MockedDbLedgerStorage extends DbLedgerStorage { @Override protected SingleDirectoryDbLedgerStorage newSingleDirectoryDbLedgerStorage(ServerConfiguration conf, - LedgerManager ledgerManager, LedgerDirsManager ledgerDirsManager, LedgerDirsManager indexDirsManager, - EntryLogger entryLogger, StatsLogger statsLogger, - long writeCacheSize, long readCacheSize, int readAheadCacheBatchSize, long readAheadCacheBatchBytesSize) + LedgerManager ledgerManager, LedgerManagerFactory ledgerManagerFactory, LedgerDirsManager ledgerDirsManager, + LedgerDirsManager indexDirsManager, EntryLogger entryLogger, StatsLogger statsLogger, long writeCacheSize, + long readCacheSize, int readAheadCacheBatchSize, long readAheadCacheBatchBytesSize) throws IOException { - return new MockedSingleDirectoryDbLedgerStorage(conf, ledgerManager, ledgerDirsManager, indexDirsManager, - entryLogger, statsLogger, allocator, writeCacheSize, + return new MockedSingleDirectoryDbLedgerStorage(conf, ledgerManager, ledgerManagerFactory, + ledgerDirsManager, indexDirsManager, entryLogger, statsLogger, allocator, writeCacheSize, readCacheSize, readAheadCacheBatchSize, readAheadCacheBatchBytesSize); } private static class MockedSingleDirectoryDbLedgerStorage extends SingleDirectoryDbLedgerStorage { public MockedSingleDirectoryDbLedgerStorage(ServerConfiguration conf, LedgerManager ledgerManager, - LedgerDirsManager ledgerDirsManager, LedgerDirsManager indexDirsManager, EntryLogger entryLogger, - StatsLogger statsLogger, + LedgerManagerFactory ledgerManagerFactory, LedgerDirsManager ledgerDirsManager, + LedgerDirsManager indexDirsManager, EntryLogger entryLogger, StatsLogger statsLogger, ByteBufAllocator allocator, long writeCacheSize, long readCacheSize, int readAheadCacheBatchSize, long readAheadCacheBatchBytesSize) throws IOException { - super(conf, ledgerManager, ledgerDirsManager, indexDirsManager, entryLogger, + super(conf, ledgerManager, ledgerManagerFactory, ledgerDirsManager, indexDirsManager, entryLogger, statsLogger, allocator, writeCacheSize, readCacheSize, readAheadCacheBatchSize, readAheadCacheBatchBytesSize); } diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/test/BookKeeperClusterTestCase.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/test/BookKeeperClusterTestCase.java index b73a3ee7b44..0cd2bbe2200 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/test/BookKeeperClusterTestCase.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/test/BookKeeperClusterTestCase.java @@ -886,7 +886,7 @@ public ServerTester(ServerConfiguration conf) throws Exception { UncleanShutdownDetection uncleanShutdownDetection = new UncleanShutdownDetectionImpl(ledgerDirsManager); storage = BookieResources.createLedgerStorage( - conf, ledgerManager, ledgerDirsManager, indexDirsManager, + conf, ledgerManager, ledgerManagerFactory, ledgerDirsManager, indexDirsManager, bookieStats, allocator); if (conf.isForceReadOnlyBookie()) { From a6e73f0db060b1011cf29f9e024d44a7dd01d26a Mon Sep 17 00:00:00 2001 From: gaozhangmin Date: Mon, 13 Jul 2026 19:24:50 +0800 Subject: [PATCH 5/5] fix error --- .../org/apache/bookkeeper/test/BookKeeperClusterTestCase.java | 2 +- .../apache/bookkeeper/stream/server/service/BookieService.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/bookkeeper-server/src/test/java/org/apache/bookkeeper/test/BookKeeperClusterTestCase.java b/bookkeeper-server/src/test/java/org/apache/bookkeeper/test/BookKeeperClusterTestCase.java index 0cd2bbe2200..3e39eaee7bc 100644 --- a/bookkeeper-server/src/test/java/org/apache/bookkeeper/test/BookKeeperClusterTestCase.java +++ b/bookkeeper-server/src/test/java/org/apache/bookkeeper/test/BookKeeperClusterTestCase.java @@ -886,7 +886,7 @@ public ServerTester(ServerConfiguration conf) throws Exception { UncleanShutdownDetection uncleanShutdownDetection = new UncleanShutdownDetectionImpl(ledgerDirsManager); storage = BookieResources.createLedgerStorage( - conf, ledgerManager, ledgerManagerFactory, ledgerDirsManager, indexDirsManager, + conf, ledgerManager, lmFactory, ledgerDirsManager, indexDirsManager, bookieStats, allocator); if (conf.isForceReadOnlyBookie()) { diff --git a/stream/server/src/main/java/org/apache/bookkeeper/stream/server/service/BookieService.java b/stream/server/src/main/java/org/apache/bookkeeper/stream/server/service/BookieService.java index d307ccbfd5b..4f496e00d5e 100644 --- a/stream/server/src/main/java/org/apache/bookkeeper/stream/server/service/BookieService.java +++ b/stream/server/src/main/java/org/apache/bookkeeper/stream/server/service/BookieService.java @@ -100,7 +100,7 @@ public BookieService(BookieConfiguration conf, StatsLogger statsLogger, LedgerDirsManager indexDirsManager = BookieResources.createIndexDirsManager( serverConf, diskChecker, bookieStats.scope(LD_INDEX_SCOPE), ledgerDirsManager); LedgerStorage storage = BookieResources.createLedgerStorage( - serverConf, ledgerManager, ledgerDirsManager, indexDirsManager, bookieStats, allocator); + serverConf, ledgerManager, lmFactory, ledgerDirsManager, indexDirsManager, bookieStats, allocator); UncleanShutdownDetection uncleanShutdownDetection = new UncleanShutdownDetectionImpl(ledgerDirsManager); LegacyCookieValidation cookieValidation = new LegacyCookieValidation(serverConf, rm);