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);