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 104e6ccc0ef..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, 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 441627737a5..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,16 +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 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(); @@ -235,16 +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()); - // instantiate zookeeper client to initialize ledger manager - - @Cleanup - MetadataBookieDriver metadataDriver = instantiateMetadataDriver(conf, statsLogger); - - @Cleanup - LedgerManagerFactory lmf = metadataDriver.getLedgerManagerFactory(); + 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 { @@ -324,22 +322,6 @@ private Set removeOverReplicatedledgers(Set bkActiveledgers, final G return overReplicatedLedgers; } - 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); - } - } - 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 ad2f0c95480..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,7 +35,9 @@ import static org.junit.Assert.assertTrue; import static org.mockito.Mockito.any; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; import static org.mockito.MockitoAnnotations.openMocks; @@ -48,6 +50,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.meta.MockLedgerManager; import org.apache.bookkeeper.stats.NullStatsLogger; import org.apache.bookkeeper.stats.StatsLogger; @@ -145,6 +148,19 @@ public void testCalculateUsageBucket() { Assert.assertEquals("Incorrect number of items", items + 1, sum); } + @Test + public void testGarbageCollectorThreadDoesNotCloseServerOwnedLedgerManagerFactory() throws Exception { + LedgerManagerFactory lmf = mock(LedgerManagerFactory.class); + File ledgerDir = tmpDirs.createNew("testGcFactoryOwnership", "ledgers"); + GarbageCollectorThread gcThread = new GarbageCollectorThread( + TestBKConfiguration.newServerConfiguration(), new MockLedgerManager(), lmf, newDirsManager(ledgerDir), + new MockLedgerStorage(), newLegacyEntryLogger(20000, ledgerDir), NullStatsLogger.INSTANCE); + + gcThread.shutdown(); + + verify(lmf, never()).close(); + } + @Test public void testExtractMetaFromEntryLogsLegacy() throws Exception { File ledgerDir = tmpDirs.createNew("testExtractMeta", "ledgers"); 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..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, 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);