Skip to content
Closed
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
1 change: 1 addition & 0 deletions fdbclient/ServerKnobs.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -794,6 +794,7 @@ void ServerKnobs::initialize(Randomize randomize, ClientKnobs* clientKnobs, IsSi
init( BACKUP_FILE_BLOCK_BYTES, 1024 * 1024 );
init( BACKUP_WORKER_LOCK_BYTES, 3e9 ); if(randomize && BUGGIFY) BACKUP_WORKER_LOCK_BYTES = deterministicRandom()->randomInt(2048, 4096) * 4096;
init( BACKUP_UPLOAD_DELAY, 10.0 ); if(randomize && BUGGIFY) BACKUP_UPLOAD_DELAY = deterministicRandom()->random01() * 60;
init( CC_RERECRUIT_BACKUP_WORKER_ENABLED, true );

//Cluster Controller
init( CLUSTER_CONTROLLER_LOGGING_DELAY, 5.0 );
Expand Down
1 change: 1 addition & 0 deletions fdbclient/include/fdbclient/ServerKnobs.h
Original file line number Diff line number Diff line change
Expand Up @@ -757,6 +757,7 @@ class SWIFT_CXX_IMMORTAL_SINGLETON_TYPE ServerKnobs : public KnobsImpl<ServerKno
int BACKUP_FILE_BLOCK_BYTES;
int64_t BACKUP_WORKER_LOCK_BYTES;
double BACKUP_UPLOAD_DELAY;
bool CC_RERECRUIT_BACKUP_WORKER_ENABLED; // re-recruit an old epoch's failed backup worker

// Cluster Controller
double CLUSTER_CONTROLLER_LOGGING_DELAY;
Expand Down
94 changes: 94 additions & 0 deletions fdbserver/ClusterRecovery.actor.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -665,6 +665,90 @@ ACTOR static Future<Optional<Version>> getMinBackupVersion(Reference<ClusterReco
}
}

// Re-recruits an old epoch's backup worker when it dies before finishing its range. Without this the
// dead worker's slot pins oldestBackupEpoch forever, which defers TLog pops and keeps any
// partitioned-log backup from becoming restorable.
//
// Errors are retried, never thrown: escaping here would fail cluster recovery, which is what
// re-recruiting outside of recovery exists to avoid.
ACTOR static Future<Void> monitorOldEpochBackupWorker(Reference<ClusterRecoveryData> self,
Database cx,
InitializeBackupRequest req,
BackupInterface interf) {
state int nextWorker = 0;
loop {
wait(waitFailureClient(interf.waitFailure,
SERVER_KNOBS->BACKUP_TIMEOUT,
-SERVER_KNOBS->BACKUP_TIMEOUT / SERVER_KNOBS->SECONDS_BEFORE_NO_FAILURE_DELAY,
/*trace=*/true));

// A worker that completed its range also stops answering, so only durable progress can say
// whether work remains.
state Reference<BackupProgress> progress(
new BackupProgress(self->dbgid, self->logSystem->getOldEpochTagsVersionsInfo()));
wait(getBackupProgress(cx, self->dbgid, progress, /*logging=*/false));

std::map<Tag, Version> status = progress->getEpochStatus(req.backupEpoch);
auto it = status.find(req.routerTag);
Version savedVersion = it == status.end() ? invalidVersion : it->second;

if (savedVersion >= req.endVersion.get()) {
// Durable progress covers the whole range, so those mutations reached the container and
// the slot can be released even though the worker died before reporting done.
CODE_PROBE(true, "Released an old epoch backup worker slot from durable progress");
TraceEvent("BackupWorkerReplacementDone", self->dbgid)
.detail("BackupEpoch", req.backupEpoch)
.detail("Tag", req.routerTag.toString())
.detail("WorkerID", interf.id())
.detail("SavedVersion", savedVersion)
.detail("EndVersion", req.endVersion.get());
self->logSystem->releaseBackupWorker(interf.id(), req.backupEpoch);
return Void();
}

InitializeBackupRequest replacement(deterministicRandom()->randomUniqueID());
replacement.recruitedEpoch = req.recruitedEpoch;
replacement.backupEpoch = req.backupEpoch;
replacement.routerTag = req.routerTag;
replacement.totalTags = req.totalTags;
replacement.startVersion = std::max(req.startVersion, savedVersion + 1);
replacement.endVersion = req.endVersion;

WorkerInterface worker = self->backupWorkers[nextWorker++ % self->backupWorkers.size()];
CODE_PROBE(true, "Re-recruiting an old epoch backup worker that died before finishing");
TraceEvent("BackupWorkerReplacement", self->dbgid)
.detail("RequestID", replacement.reqId)
.detail("Tag", replacement.routerTag.toString())
.detail("Epoch", replacement.recruitedEpoch)
.detail("BackupEpoch", replacement.backupEpoch)
.detail("StartVersion", replacement.startVersion)
.detail("EndVersion", replacement.endVersion.get())
.detail("DeadWorkerID", interf.id());

try {
InitializeBackupReply reply = wait(throwErrorOr(
worker.backup.getReplyUnlessFailedFor(replacement,
SERVER_KNOBS->BACKUP_TIMEOUT,
SERVER_KNOBS->MASTER_FAILURE_SLOPE_DURING_RECOVERY)));
if (!self->logSystem->replaceBackupWorker(interf.id(), reply)) {
return Void();
}
interf = reply.interf;
} catch (Error& e) {
if (e.code() == error_code_actor_cancelled) {
throw;
}
CODE_PROBE(true, "Old epoch backup worker re-recruitment failed", probe::decoration::rare);
TraceEvent(SevWarn, "BackupWorkerReplacementFailed", self->dbgid)
.error(e)
.detail("BackupEpoch", req.backupEpoch)
.detail("Tag", req.routerTag.toString())
.detail("DeadWorkerID", interf.id());
wait(delay(SERVER_KNOBS->BACKUP_TIMEOUT));
}
}
}

ACTOR static Future<Void> recruitBackupWorkers(Reference<ClusterRecoveryData> self, Database cx) {
ASSERT(self->backupWorkers.size() > 0);

Expand All @@ -676,6 +760,7 @@ ACTOR static Future<Void> recruitBackupWorkers(Reference<ClusterRecoveryData> se
new BackupProgress(self->dbgid, self->logSystem->getOldEpochTagsVersionsInfo()));
state Future<Void> gotProgress = getBackupProgress(cx, self->dbgid, backupProgress, /*logging=*/true);
state std::vector<Future<InitializeBackupReply>> initializationReplies;
state std::vector<InitializeBackupRequest> oldEpochRequests; // requests for previous epochs' unfinished work

state std::vector<std::pair<UID, Tag>> idsTags; // worker IDs and tags for current epoch
state int logRouterTags = self->logSystem->getLogRouterTags();
Expand Down Expand Up @@ -744,13 +829,22 @@ ACTOR static Future<Void> recruitBackupWorkers(Reference<ClusterRecoveryData> se
throwErrorOr(worker.backup.getReplyUnlessFailedFor(
req, SERVER_KNOBS->BACKUP_TIMEOUT, SERVER_KNOBS->MASTER_FAILURE_SLOPE_DURING_RECOVERY)),
backup_worker_failed()));
oldEpochRequests.push_back(req);
}
}

std::vector<InitializeBackupReply> newRecruits = wait(getAll(initializationReplies));
self->logSystem->setBackupWorkers(newRecruits);
TraceEvent("BackupRecruitmentDone", self->dbgid).log();
self->registrationTrigger.trigger();

ASSERT(newRecruits.size() == logRouterTags + oldEpochRequests.size());
if (SERVER_KNOBS->CC_RERECRUIT_BACKUP_WORKER_ENABLED) {
for (int j = 0; j < oldEpochRequests.size(); j++) {
self->addActor.send(monitorOldEpochBackupWorker(
self, cx, oldEpochRequests[j], newRecruits[logRouterTags + j].interf));
}
}
return Void();
}

Expand Down
69 changes: 62 additions & 7 deletions fdbserver/TagPartitionedLogSystem.actor.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -1939,13 +1939,7 @@ bool TagPartitionedLogSystem::removeBackupWorker(const BackupWorkerDoneRequest&
}

if (removed) {
oldestBackupEpoch = epoch;
for (const auto& old : oldLogData) {
if (old.epoch < oldestBackupEpoch && old.tLogs[0]->backupWorkers.size() > 0) {
oldestBackupEpoch = old.epoch;
}
}
backupWorkerChanged.trigger();
recomputeOldestBackupEpoch();
} else {
removedBackupWorkers.insert(req.workerUID);
}
Expand All @@ -1958,6 +1952,67 @@ bool TagPartitionedLogSystem::removeBackupWorker(const BackupWorkerDoneRequest&
return removed;
}

bool TagPartitionedLogSystem::replaceBackupWorker(UID deadWorker, const InitializeBackupReply& reply) {
Reference<LogSet> logset = getEpochLogSet(reply.backupEpoch);
if (!logset.isValid()) {
return false;
}

for (auto& worker : logset->backupWorkers) {
if (worker->get().interf().id() != deadWorker) {
continue;
}
// Keeping the entry in place leaves the epoch's count unchanged, so it retains its hold on
// oldestBackupEpoch and no TLog data becomes collectable while its work is outstanding.
worker->setUnconditional(OptionalInterface<BackupInterface>(reply.interf));
TraceEvent("ReplaceBackupWorker", dbgid)
.detail("BackupEpoch", reply.backupEpoch)
.detail("DeadWorkerID", deadWorker)
.detail("WorkerID", reply.interf.id());

// A replacement that finished before this call found no entry to erase and only recorded its
// UID. Honour that now, or the epoch keeps a slot for a worker that will never report again.
if (removedBackupWorkers.contains(reply.interf.id())) {
removedBackupWorkers.erase(reply.interf.id());
releaseBackupWorker(reply.interf.id(), reply.backupEpoch);
return false;
}

backupWorkerChanged.trigger();
return true;
}
return false;
}

void TagPartitionedLogSystem::releaseBackupWorker(UID worker, LogEpoch backupEpoch) {
Reference<LogSet> logset = getEpochLogSet(backupEpoch);
if (!logset.isValid()) {
return;
}

for (auto it = logset->backupWorkers.begin(); it != logset->backupWorkers.end(); it++) {
if (it->getPtr()->get().interf().id() == worker) {
logset->backupWorkers.erase(it);
recomputeOldestBackupEpoch();
TraceEvent("ReleaseBackupWorker", dbgid)
.detail("BackupEpoch", backupEpoch)
.detail("WorkerID", worker)
.detail("OldestBackupEpoch", oldestBackupEpoch);
return;
}
}
}

void TagPartitionedLogSystem::recomputeOldestBackupEpoch() {
oldestBackupEpoch = epoch;
for (const auto& old : oldLogData) {
if (old.epoch < oldestBackupEpoch && old.tLogs[0]->backupWorkers.size() > 0) {
oldestBackupEpoch = old.epoch;
}
}
backupWorkerChanged.trigger();
}

LogEpoch TagPartitionedLogSystem::getOldestBackupEpoch() const {
return oldestBackupEpoch;
}
Expand Down
10 changes: 10 additions & 0 deletions fdbserver/include/fdbserver/LogSystem.h
Original file line number Diff line number Diff line change
Expand Up @@ -718,6 +718,16 @@ struct ILogSystem {
// if the worker is not found.
virtual bool removeBackupWorker(const BackupWorkerDoneRequest& req) = 0;

// Points an old epoch's backup worker slot at a re-recruited replacement, keeping the entry in
// place so the epoch retains its hold on oldestBackupEpoch. Returns false if the dead worker is
// no longer tracked, meaning the replacement must not be installed.
virtual bool replaceBackupWorker(UID deadWorker, const InitializeBackupReply& reply) = 0;

// Frees a backup worker slot whose work durable progress already shows as complete, for a worker
// that died without reporting done. Unlike removeBackupWorker this records nothing for a later
// setBackupWorkers, since there is no done request in flight to reconcile with.
virtual void releaseBackupWorker(UID worker, LogEpoch backupEpoch) = 0;

virtual LogEpoch getOldestBackupEpoch() const = 0;
virtual void setOldestBackupEpoch(LogEpoch epoch) = 0;
};
Expand Down
7 changes: 7 additions & 0 deletions fdbserver/include/fdbserver/TagPartitionedLogSystem.actor.h
Original file line number Diff line number Diff line change
Expand Up @@ -344,6 +344,13 @@ struct TagPartitionedLogSystem final : ILogSystem, ReferenceCounted<TagPartition

bool removeBackupWorker(const BackupWorkerDoneRequest& req) final;

bool replaceBackupWorker(UID deadWorker, const InitializeBackupReply& reply) final;

void releaseBackupWorker(UID worker, LogEpoch backupEpoch) final;

// Pins oldestBackupEpoch to the lowest old epoch still holding backup workers.
void recomputeOldestBackupEpoch();

LogEpoch getOldestBackupEpoch() const final;

void setOldestBackupEpoch(LogEpoch epoch) final;
Expand Down
Loading