Merge pull request #5787 from dlambrig/integrate-PR5700

version vector / Calculate TPCV on resolvers
This commit is contained in:
sbodagala 2021-10-20 16:35:44 -04:00 committed by GitHub
commit 48a0ecd647
No known key found for this signature in database
GPG Key ID: 4AEE18F83AFDEB23
7 changed files with 56 additions and 137 deletions

View File

@ -37,7 +37,7 @@ void ServerKnobs::initialize(Randomize randomize, ClientKnobs* clientKnobs, IsSi
init( MAX_WRITE_TRANSACTION_LIFE_VERSIONS, 5 * VERSIONS_PER_SECOND ); if (randomize && BUGGIFY) MAX_WRITE_TRANSACTION_LIFE_VERSIONS=std::max<int>(1, 1 * VERSIONS_PER_SECOND);
init( MAX_COMMIT_BATCH_INTERVAL, 2.0 ); if( randomize && BUGGIFY ) MAX_COMMIT_BATCH_INTERVAL = 0.5; // Each commit proxy generates a CommitTransactionBatchRequest at least this often, so that versions always advance smoothly
MAX_COMMIT_BATCH_INTERVAL = std::min(MAX_COMMIT_BATCH_INTERVAL, MAX_READ_TRANSACTION_LIFE_VERSIONS/double(2*VERSIONS_PER_SECOND)); // Ensure that the proxy commits 2 times every MAX_READ_TRANSACTION_LIFE_VERSIONS, otherwise the master will not give out versions fast enough
init( ENABLE_VERSION_VECTOR, true );
init( ENABLE_VERSION_VECTOR, true );
init( ENABLE_VERSION_VECTOR_TLOG_UNICAST, false );
// TLogs

View File

@ -536,12 +536,9 @@ struct CommitBatchContext {
double commitStartTime;
std::unordered_map<uint16_t, Version> tpcvMap; // obtained from sequencer
std::set<uint16_t> writtenTLogs; // the set of tlog locations written to in the mutation.
std::unordered_map<uint16_t, Version> tpcvMap; // obtained from resolver
std::set<Tag> writtenTags; // final set tags written to in the batch
std::set<Tag> writtenTagsPreResolution; // tags written to in the batch not including any changes from the resolver.
bool hasMetadataMutation = false;
bool metadataMutationFromProxy = false;
CommitBatchContext(ProxyCommitData*, const std::vector<CommitTransactionRequest>*, const int);
@ -564,9 +561,7 @@ std::set<Tag> CommitBatchContext::getWrittenTagsPreResolution() {
if (isSingleKeyMutation((MutationRef::Type)m.type)) {
auto& tags = pProxyCommitData->tagsForKey(m.param1);
transactionTags.insert(tags.begin(), tags.end());
toCommit.getLocations(tags, writtenTLogs);
if (pProxyCommitData->cacheInfo[m.param1]) {
toCommit.getLocations(cacheVector, writtenTLogs);
transactionTags.insert(cacheTag);
}
} else if (m.type == MutationRef::ClearRange) {
@ -579,7 +574,6 @@ std::set<Tag> CommitBatchContext::getWrittenTagsPreResolution() {
ranges.begin().value().populateTags();
filteredTags.insert(ranges.begin().value().tags.begin(), ranges.begin().value().tags.end());
transactionTags.insert(ranges.begin().value().tags.begin(), ranges.begin().value().tags.end());
toCommit.getLocations(filteredTags, writtenTLogs);
} else {
std::set<Tag> allSources;
for (auto r : ranges) {
@ -587,17 +581,14 @@ std::set<Tag> CommitBatchContext::getWrittenTagsPreResolution() {
allSources.insert(r.value().tags.begin(), r.value().tags.end());
transactionTags.insert(r.value().tags.begin(), r.value().tags.end());
}
toCommit.getLocations(allSources, writtenTLogs);
}
if (pProxyCommitData->needsCacheTag(clearRange)) {
toCommit.getLocations(cacheVector, writtenTLogs);
transactionTags.insert(cacheTag);
}
} else {
UNREACHABLE();
}
}
hasMetadataMutation = containsMetadataMutation(trs[transactionNum].transaction.mutations);
}
return transactionTags;
@ -733,26 +724,15 @@ ACTOR Future<Void> preresolutionProcessing(CommitBatchContext* self) {
if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) {
self->writtenTagsPreResolution = self->getWrittenTagsPreResolution();
if (self->hasMetadataMutation) {
int numLogs = pProxyCommitData->db->get().logSystemConfig.numLogs();
for (int i = 0; i < numLogs; i++) {
self->writtenTLogs.insert(i);
}
}
}
GetCommitVersionRequest req(span.context,
pProxyCommitData->commitVersionRequestNumber++,
pProxyCommitData->mostRecentProcessedRequestNumber,
pProxyCommitData->dbgid,
self->writtenTLogs);
pProxyCommitData->dbgid);
state double beforeGettingCommitVersion = now();
GetCommitVersionReply versionReply = wait(brokenPromiseToNever(
pProxyCommitData->master.getCommitVersion.getReply(req, TaskPriority::ProxyMasterVersionReply)));
if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) {
self->tpcvMap = versionReply.tpcvMap;
}
pProxyCommitData->mostRecentProcessedRequestNumber = versionReply.requestNum;
pProxyCommitData->stats.txnCommitVersionAssigned += trs.size();
@ -807,6 +787,7 @@ ACTOR Future<Void> getResolution(CommitBatchContext* self) {
std::vector<Future<ResolveTransactionBatchReply>> replies;
for (int r = 0; r < pProxyCommitData->resolvers.size(); r++) {
requests.requests[r].debugID = self->debugID;
requests.requests[r].writtenTags = self->writtenTagsPreResolution;
replies.push_back(trackResolutionMetrics(pProxyCommitData->stats.resolverDist[r],
brokenPromiseToNever(pProxyCommitData->resolvers[r].resolve.getReply(
requests.requests[r], TaskPriority::ProxyResolverReply))));
@ -848,6 +829,9 @@ void assertResolutionStateMutationsSizeConsistent(const std::vector<ResolveTrans
for (int r = 1; r < resolution.size(); r++) {
ASSERT(resolution[r].stateMutations.size() == resolution[0].stateMutations.size());
ASSERT_EQ(resolution[0].privateMutationCount, resolution[r].privateMutationCount);
if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) {
ASSERT_EQ(resolution[0].tpcvMap.size(), resolution[r].tpcvMap.size());
}
for (int s = 0; s < resolution[r].stateMutations.size(); s++) {
ASSERT(resolution[r].stateMutations[s].size() == resolution[0].stateMutations[s].size());
}
@ -881,11 +865,6 @@ void applyMetadataEffect(CommitBatchContext* self) {
self->forceRecovery,
/* popVersion= */ 0,
/* initialCommit */ false);
if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST &&
containsMetadataMutation(
self->resolution[0].stateMutations[versionIndex][transactionIndex].mutations)) {
self->metadataMutationFromProxy = true;
}
}
if (self->resolution[0].stateMutations[versionIndex][transactionIndex].mutations.size() &&
self->firstStateMutations) {
@ -994,6 +973,10 @@ ACTOR Future<Void> applyMetadataToCommittedTransactions(CommitBatchContext* self
// in the same set of transactions as this proxy.
ResolveTransactionBatchReply& reply = self->resolution[0];
self->toCommit.setMutations(reply.privateMutationCount, reply.privateMutations);
if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) {
// TraceEvent("ResolverReturn").detail("ReturnTags",reply.writtenTags).detail("TPCVsize",reply.tpcvMap.size()).detail("ReqTags",self->writtenTagsPreResolution);
self->tpcvMap = reply.tpcvMap;
}
}
self->lockedKey = pProxyCommitData->txnStateStore->readValue(databaseLockedKey).get();
@ -1018,16 +1001,6 @@ ACTOR Future<Void> applyMetadataToCommittedTransactions(CommitBatchContext* self
return Void();
}
ACTOR Future<Void> getTPCV(CommitBatchContext* self, std::set<uint16_t> writtenTLogs) {
state ProxyCommitData* const pProxyCommitData = self->pProxyCommitData;
GetTLogPrevCommitVersionReply rep =
wait(brokenPromiseToNever(pProxyCommitData->master.getTLogPrevCommitVersion.getReply(
GetTLogPrevCommitVersionRequest(writtenTLogs, self->commitVersion, self->prevVersion))));
// TraceEvent("GetTLogPrevCommitVersionRequest");
self->tpcvMap.insert(rep.tpcvMap.begin(), rep.tpcvMap.end());
return Void();
}
/// This second pass through committed transactions assigns the actual mutations to the appropriate storage servers'
/// tags
ACTOR Future<Void> assignMutationsToStorageServers(CommitBatchContext* self) {
@ -1242,14 +1215,6 @@ ACTOR Future<Void> postResolution(CommitBatchContext* self) {
wait(Future<Void>(Never()));
}
if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST && self->metadataMutationFromProxy &&
!self->hasMetadataMutation) {
// TraceEvent("Abort metadataMutationFromProxy");
for (int transactionNum = 0; transactionNum < trs.size(); transactionNum++) {
self->committed[transactionNum] = ConflictBatch::TransactionConflict;
}
}
// First pass
wait(applyMetadataToCommittedTransactions(self));
@ -1258,25 +1223,6 @@ ACTOR Future<Void> postResolution(CommitBatchContext* self) {
self->toCommit.saveTags(self->writtenTags);
if (self->writtenTags.size() && !self->metadataMutationFromProxy &&
SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) {
// confirm all serialized tags are sent to a tLog for which the previous commit version was obtained.
std::set<uint16_t> postResolutionTLogs;
self->toCommit.getLocations(self->writtenTags, postResolutionTLogs);
for (auto& t : postResolutionTLogs) {
if (self->writtenTLogs.find(t) == self->writtenTLogs.end()) {
TraceEvent(SevError, "TagHasNoPCV", pProxyCommitData->dbgid)
.detail("tagsBeforeResolution", self->writtenTagsPreResolution)
.detail("tagsAfterResolution", self->writtenTags)
.detail("numTLogsBeforeResolution", self->writtenTLogs.size())
.detail("numTLogsAfterResolution", postResolutionTLogs.size())
.detail("hasMetadataMutation", self->hasMetadataMutation)
.detail("metadataMutationFromProxy", self->metadataMutationFromProxy);
ASSERT(false);
}
}
}
// Serialize and backup the mutations as a single mutation
if ((pProxyCommitData->vecBackupKeys.size() > 1) && self->logRangeMutations.size()) {
wait(addBackupMutations(pProxyCommitData,

View File

@ -44,7 +44,6 @@ struct MasterInterface {
RequestStream<struct GetRawCommittedVersionRequest> getLiveCommittedVersion;
// Report a proxy's committed version.
RequestStream<struct ReportRawCommittedVersionRequest> reportLiveCommittedVersion;
RequestStream<struct GetTLogPrevCommitVersionRequest> getTLogPrevCommitVersion;
NetworkAddress address() const { return changeCoordinators.getEndpoint().getPrimaryAddress(); }
NetworkAddressList addresses() const { return changeCoordinators.getEndpoint().addresses; }
@ -68,8 +67,6 @@ struct MasterInterface {
RequestStream<struct GetRawCommittedVersionRequest>(waitFailure.getEndpoint().getAdjustedEndpoint(5));
reportLiveCommittedVersion = RequestStream<struct ReportRawCommittedVersionRequest>(
waitFailure.getEndpoint().getAdjustedEndpoint(6));
getTLogPrevCommitVersion =
RequestStream<struct GetTLogPrevCommitVersionRequest>(waitFailure.getEndpoint().getAdjustedEndpoint(7));
}
}
@ -82,7 +79,6 @@ struct MasterInterface {
streams.push_back(notifyBackupWorkerDone.getReceiver());
streams.push_back(getLiveCommittedVersion.getReceiver(TaskPriority::GetLiveCommittedVersion));
streams.push_back(reportLiveCommittedVersion.getReceiver(TaskPriority::ReportLiveCommittedVersion));
streams.push_back(getTLogPrevCommitVersion.getReceiver(TaskPriority::GetTLogPrevCommitVersion));
FlowTransport::transport().addEndpoints(streams);
}
};
@ -156,7 +152,6 @@ struct GetCommitVersionReply {
Version version;
Version prevVersion;
uint64_t requestNum;
std::unordered_map<uint16_t, Version> tpcvMap;
GetCommitVersionReply() : resolverChangesVersion(0), version(0), prevVersion(0), requestNum(0) {}
explicit GetCommitVersionReply(Version version, Version prevVersion, uint64_t requestNum)
@ -164,7 +159,7 @@ struct GetCommitVersionReply {
template <class Ar>
void serialize(Ar& ar) {
serializer(ar, resolverChanges, resolverChangesVersion, version, prevVersion, requestNum, tpcvMap);
serializer(ar, resolverChanges, resolverChangesVersion, version, prevVersion, requestNum);
}
};
@ -174,46 +169,28 @@ struct GetCommitVersionRequest {
uint64_t requestNum;
uint64_t mostRecentProcessedRequestNum;
UID requestingProxy;
std::set<uint16_t> writtenTLogs;
ReplyPromise<GetCommitVersionReply> reply;
GetCommitVersionRequest() {}
GetCommitVersionRequest(SpanID spanContext,
uint64_t requestNum,
uint64_t mostRecentProcessedRequestNum,
UID requestingProxy,
std::set<uint16_t>& writtenTLogs)
UID requestingProxy)
: spanContext(spanContext), requestNum(requestNum), mostRecentProcessedRequestNum(mostRecentProcessedRequestNum),
requestingProxy(requestingProxy), writtenTLogs(writtenTLogs) {}
requestingProxy(requestingProxy) {}
template <class Ar>
void serialize(Ar& ar) {
serializer(ar, requestNum, mostRecentProcessedRequestNum, requestingProxy, writtenTLogs, reply, spanContext);
serializer(ar, requestNum, mostRecentProcessedRequestNum, requestingProxy, reply, spanContext);
}
};
struct GetTLogPrevCommitVersionReply {
constexpr static FileIdentifier file_identifier = 16683183;
std::unordered_map<uint16_t, Version> tpcvMap;
GetTLogPrevCommitVersionReply() {}
template <class Ar>
void serialize(Ar& ar) {
serializer(ar, tpcvMap);
}
};
struct GetTLogPrevCommitVersionRequest {
constexpr static FileIdentifier file_identifier = 16683184;
std::set<uint16_t> writtenTLogs;
Version commitVersion;
Version prev;
ReplyPromise<GetTLogPrevCommitVersionReply> reply;
GetTLogPrevCommitVersionRequest() {}
GetTLogPrevCommitVersionRequest(std::set<uint16_t>& writtenTLogs, Version commitVersion, Version prev)
: writtenTLogs(writtenTLogs), commitVersion(commitVersion), prev(prev) {}
template <class Ar>
void serialize(Ar& ar) {
serializer(ar, writtenTLogs, commitVersion, prev, reply);
serializer(ar);
}
};

View File

@ -78,6 +78,9 @@ struct Resolver : ReferenceCounted<Resolver> {
Version debugMinRecentStateVersion = 0;
// The previous commit versions per tlog
std::vector<Version> tpcvVector;
CounterCollection cc;
Counter resolveBatchIn;
Counter resolveBatchStart;
@ -94,6 +97,7 @@ struct Resolver : ReferenceCounted<Resolver> {
Counter resolveBatchOut;
Counter metricsRequests;
Counter splitRequests;
int numLogs;
Future<Void> logger;
@ -192,6 +196,7 @@ ACTOR Future<Void> resolveBatch(Reference<Resolver> self, ResolveTransactionBatc
g_traceBatch.addEvent("CommitDebug", debugID.get().first(), "Resolver.resolveBatch.AfterOrderer");
ResolveTransactionBatchReply& reply = proxyInfo.outstandingBatches[req.version];
reply.writtenTags = req.writtenTags;
std::vector<int> commitList;
std::vector<int> tooOldList;
@ -281,6 +286,10 @@ ACTOR Future<Void> resolveBatch(Reference<Resolver> self, ResolveTransactionBatc
reply.privateMutations.push_back(reply.arena, mutations);
reply.arena.dependsOn(mutations.arena());
}
if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) {
// merge mutation tags with sent client tags
toCommit.saveTags(reply.writtenTags);
}
reply.privateMutationCount = toCommit.getMutationCount();
}
@ -340,6 +349,23 @@ ACTOR Future<Void> resolveBatch(Reference<Resolver> self, ResolveTransactionBatc
}
}
if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) {
std::set<uint16_t> writtenTLogs;
if (reply.privateMutationCount) {
for (int i = 0; i < self->numLogs; i++) {
writtenTLogs.insert(i);
}
} else {
toCommit.getLocations(reply.writtenTags, writtenTLogs);
}
if (self->tpcvVector[0] == invalidVersion) {
std::fill(self->tpcvVector.begin(), self->tpcvVector.end(), req.prevVersion);
}
for (uint16_t tLog : writtenTLogs) {
reply.tpcvMap[tLog] = self->tpcvVector[tLog];
self->tpcvVector[tLog] = req.version;
}
}
self->version.set(req.version);
bool breachedLimit = self->totalStateBytes.get() <= SERVER_KNOBS->RESOLVER_STATE_MEMORY_LIMIT &&
self->totalStateBytes.get() + stateBytes > SERVER_KNOBS->RESOLVER_STATE_MEMORY_LIMIT;
@ -568,6 +594,12 @@ ACTOR Future<Void> resolverCore(ResolverInterface resolver,
// This has to be declared after the self->txnStateStore get initialized
transactionStateResolveContext = TransactionStateResolveContext(self, &addActor);
if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) {
self->numLogs = db->get().logSystemConfig.numLogs();
self->tpcvVector.resize(1 + self->numLogs, 0);
std::fill(self->tpcvVector.begin(), self->tpcvVector.end(), invalidVersion);
}
}
loop choose {

View File

@ -96,6 +96,9 @@ struct ResolveTransactionBatchReply {
VectorRef<StringRef> privateMutations;
uint32_t privateMutationCount;
std::unordered_map<uint16_t, Version> tpcvMap;
std::set<Tag> writtenTags;
template <class Archive>
void serialize(Archive& ar) {
serializer(ar,
@ -105,6 +108,8 @@ struct ResolveTransactionBatchReply {
conflictingKeyRangeMap,
privateMutations,
privateMutationCount,
tpcvMap,
writtenTags,
arena);
}
};
@ -123,6 +128,8 @@ struct ResolveTransactionBatchRequest {
ReplyPromise<ResolveTransactionBatchReply> reply;
Optional<UID> debugID;
std::set<Tag> writtenTags;
template <class Archive>
void serialize(Archive& ar) {
serializer(ar,
@ -134,6 +141,7 @@ struct ResolveTransactionBatchRequest {
reply,
arena,
debugID,
writtenTags,
spanContext);
}
};

View File

@ -2148,9 +2148,7 @@ ACTOR Future<Void> tLogCommit(TLogData* self,
}
logData->minKnownCommittedVersion = std::max(logData->minKnownCommittedVersion, req.minKnownCommittedVersion);
wait(logData->version.whenAtLeast(req.prevVersion));
// Calling check_yield instead of yield to avoid a destruction ordering problem in simulation
if (g_network->check_yield(g_network->getCurrentTask())) {
wait(delay(0, g_network->getCurrentTask()));
@ -2198,7 +2196,6 @@ ACTOR Future<Void> tLogCommit(TLogData* self,
if (self->diskQueueCommitBytes > SERVER_KNOBS->MAX_QUEUE_COMMIT_BYTES) {
self->largeDiskQueueCommitBytes.set(true);
}
// Notifies the commitQueue actor to commit persistentQueue, and also unblocks tLogPeekMessages actors
logData->version.set(req.version);

View File

@ -250,8 +250,6 @@ struct MasterData : NonCopyable, ReferenceCounted<MasterData> {
// up-to-date in the presence of key range splits/merges.
VersionVector ssVersionVector;
// The previous commit versions per tlog
std::vector<Version> tpcvVector;
CounterCollection cc;
Counter changeCoordinatorsRequests;
Counter getCommitVersionRequests;
@ -1189,9 +1187,6 @@ ACTOR Future<Void> getVersion(Reference<MasterData> self, GetCommitVersionReques
self->lastVersionTime = now();
self->version = self->recoveryTransactionVersion;
rep.prevVersion = self->lastEpochEnd;
if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) {
std::fill(self->tpcvVector.begin(), self->tpcvVector.end(), self->lastEpochEnd);
}
} else {
double t1 = now();
@ -1228,13 +1223,6 @@ ACTOR Future<Void> getVersion(Reference<MasterData> self, GetCommitVersionReques
proxyItr->second.replies[req.requestNum] = rep;
ASSERT(rep.prevVersion >= 0);
if (SERVER_KNOBS->ENABLE_VERSION_VECTOR_TLOG_UNICAST) {
for (uint16_t tLog : req.writtenTLogs) {
rep.tpcvMap[tLog] = self->tpcvVector[tLog];
self->tpcvVector[tLog] = rep.version;
}
}
req.reply.send(rep);
ASSERT(proxyItr->second.latestRequestNum.get() == req.requestNum - 1);
@ -1282,24 +1270,6 @@ ACTOR Future<Void> waitForPrev(Reference<MasterData> self, ReportRawCommittedVer
return Void();
}
ACTOR Future<Void> waitForTLogPrev(Reference<MasterData> self, GetTLogPrevCommitVersionRequest req) {
// TraceEvent("WaitForTLogPrev").detail("Prev",req.prev).detail("CommitVersion",req.commitVersion).detail("PrevTLogVersion",self->prevTLogVersion.get());
if (self->prevTLogVersion.get() != invalidVersion) {
wait(self->prevTLogVersion.whenAtLeast(req.prev));
} else {
std::fill(self->tpcvVector.begin(), self->tpcvVector.end(), self->lastEpochEnd);
}
GetTLogPrevCommitVersionReply reply;
for (uint16_t tLog : req.writtenTLogs) {
reply.tpcvMap[tLog] = self->tpcvVector[tLog];
self->tpcvVector[tLog] = req.commitVersion;
}
self->prevTLogVersion.set(req.commitVersion);
req.reply.send(reply);
return Void();
}
ACTOR Future<Void> serveLiveCommittedVersion(Reference<MasterData> self) {
loop {
choose {
@ -1333,10 +1303,6 @@ ACTOR Future<Void> serveLiveCommittedVersion(Reference<MasterData> self) {
req.reply.send(Void());
}
}
when(GetTLogPrevCommitVersionRequest req =
waitNext(self->myInterface.getTLogPrevCommitVersion.getFuture())) {
self->addActor.send(waitForTLogPrev(self, req));
}
}
}
}
@ -1974,10 +1940,6 @@ ACTOR Future<Void> masterCore(Reference<MasterData> self) {
tr.read_snapshot = self->recoveryTransactionVersion; // lastEpochEnd would make more sense, but isn't in the initial
// window of the resolver(s)
// resize the TPCV vector to the number of tlogs and initialize to first transaction
int numLogs = self->dbInfo->get().logSystemConfig.numLogs();
self->tpcvVector.resize(1 + numLogs, 0);
TraceEvent("MasterRecoveryCommit", self->dbgid).log();
state Future<ErrorOr<CommitID>> recoveryCommit = self->commitProxies[0].commit.tryGetReply(recoveryCommitRequest);
self->addActor.send(self->logSystem->onError());
@ -2117,9 +2079,6 @@ ACTOR Future<Void> masterServer(MasterInterface mi,
wait(delay(5));
throw worker_removed();
}
// resize the TPCV vector to the number of tlogs and initialize to first transaction
int numLogs = self->dbInfo->get().logSystemConfig.numLogs();
self->tpcvVector.resize(1 + numLogs, 0);
}
when(BackupWorkerDoneRequest req = waitNext(mi.notifyBackupWorkerDone.getFuture())) {
if (self->logSystem.isValid() && self->logSystem->removeBackupWorker(req)) {