Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -363,6 +363,7 @@ public static LedgerStorage mountLedgerStorageOffline(ServerConfiguration conf,

if (null == ledgerStorage) {
ledgerStorage = BookieResources.createLedgerStorage(conf, null,
null,
ledgerDirsManager,
indexDirsManager,
statsLogger,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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,
Expand All @@ -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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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
Expand All @@ -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;

Expand All @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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;
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

/**
Expand All @@ -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);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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.
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -235,16 +236,13 @@ private Set<Long> removeOverReplicatedledgers(Set<Long> bkActiveledgers, final G
final Set<Long> 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 {
Expand Down Expand Up @@ -324,22 +322,6 @@ private Set<Long> removeOverReplicatedledgers(Set<Long> 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<LedgerMetadata> 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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -123,6 +124,7 @@ public class DbLedgerStorage implements LedgerStorage {
private static final long STORAGE_FLAGS_KEY = 0L;
private int numberOfDirs;
private List<SingleDirectoryDbLedgerStorage> ledgerStorageList;
private LedgerManagerFactory ledgerManagerFactory;

private ExecutorService entryLoggerWriteExecutor = null;
private ExecutorService entryLoggerFlushExecutor = null;
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
Loading
Loading