From 3379f1e9742e574ce5979f454110fb49c761728d Mon Sep 17 00:00:00 2001 From: Jingyu Zhou Date: Fri, 18 Mar 2022 15:14:02 -0700 Subject: [PATCH 1/6] Move resolutionBalancing() back to master This revert the behavior done by a recent refactor on master recovery in PR #6191. --- fdbserver/ClusterRecovery.actor.cpp | 151 +----------------- fdbserver/ClusterRecovery.actor.h | 16 +- fdbserver/MasterInterface.h | 19 ++- fdbserver/masterserver.actor.cpp | 188 ++++++++++++++++++++--- fdbserver/workloads/KillRegion.actor.cpp | 5 +- 5 files changed, 188 insertions(+), 191 deletions(-) diff --git a/fdbserver/ClusterRecovery.actor.cpp b/fdbserver/ClusterRecovery.actor.cpp index ac8fee9b9e..22b6807a00 100644 --- a/fdbserver/ClusterRecovery.actor.cpp +++ b/fdbserver/ClusterRecovery.actor.cpp @@ -470,7 +470,6 @@ ACTOR Future trackTlogRecovery(Reference self, self->dbgid) .detail("StatusCode", RecoveryStatus::fully_recovered) .detail("Status", RecoveryStatus::names[RecoveryStatus::fully_recovered]) - .detail("FullyRecoveredAtVersion", self->version) .detail("ClusterId", self->clusterId) .trackLatest(self->clusterRecoveryStateEventHolder->trackingKey); @@ -511,144 +510,6 @@ ACTOR Future trackTlogRecovery(Reference self, } } -std::pair findRange(CoalescedKeyRangeMap& key_resolver, - Standalone>& movedRanges, - int src, - int dest) { - auto ranges = key_resolver.ranges(); - auto prev = ranges.begin(); - auto it = ranges.begin(); - ++it; - if (it == ranges.end()) { - if (ranges.begin().value() != src || - std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(ranges.begin()->range(), dest)) != - movedRanges.end()) - throw operation_failed(); - return std::make_pair(ranges.begin().range(), true); - } - - std::set borders; - // If possible expand an existing boundary between the two resolvers - for (; it != ranges.end(); ++it) { - if (it->value() == src && prev->value() == dest && - std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(it->range(), dest)) == - movedRanges.end()) { - return std::make_pair(it->range(), true); - } - if (it->value() == dest && prev->value() == src && - std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(prev->range(), dest)) == - movedRanges.end()) { - return std::make_pair(prev->range(), false); - } - if (it->value() == dest) - borders.insert(prev->value()); - if (prev->value() == dest) - borders.insert(it->value()); - ++prev; - } - - prev = ranges.begin(); - it = ranges.begin(); - ++it; - // If possible create a new boundry which doesn't exist yet - for (; it != ranges.end(); ++it) { - if (it->value() == src && !borders.count(prev->value()) && - std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(it->range(), dest)) == - movedRanges.end()) { - return std::make_pair(it->range(), true); - } - if (prev->value() == src && !borders.count(it->value()) && - std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(prev->range(), dest)) == - movedRanges.end()) { - return std::make_pair(prev->range(), false); - } - ++prev; - } - - it = ranges.begin(); - for (; it != ranges.end(); ++it) { - if (it->value() == src && - std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(it->range(), dest)) == - movedRanges.end()) { - return std::make_pair(it->range(), true); - } - } - throw operation_failed(); // we are already attempting to move all of the data one resolver is assigned, so do not - // move anything -} - -ACTOR Future resolutionBalancing(Reference self) { - state CoalescedKeyRangeMap key_resolver; - key_resolver.insert(allKeys, 0); - loop { - wait(delay(SERVER_KNOBS->MIN_BALANCE_TIME, TaskPriority::ResolutionMetrics)); - while (self->resolverChanges.get().size()) - wait(self->resolverChanges.onChange()); - state std::vector> futures; - for (auto& p : self->resolvers) - futures.push_back( - brokenPromiseToNever(p.metrics.getReply(ResolutionMetricsRequest(), TaskPriority::ResolutionMetrics))); - wait(waitForAll(futures)); - state IndexedSet, NoMetric> metrics; - - int64_t total = 0; - for (int i = 0; i < futures.size(); i++) { - total += futures[i].get().value; - metrics.insert(std::make_pair(futures[i].get().value, i), NoMetric()); - //TraceEvent("ResolverMetric").detail("I", i).detail("Metric", futures[i].get()); - } - if (metrics.lastItem()->first - metrics.begin()->first > SERVER_KNOBS->MIN_BALANCE_DIFFERENCE) { - try { - state int src = metrics.lastItem()->second; - state int dest = metrics.begin()->second; - state int64_t amount = std::min(metrics.lastItem()->first - total / self->resolvers.size(), - total / self->resolvers.size() - metrics.begin()->first) / - 2; - state Standalone> movedRanges; - - loop { - state std::pair range = findRange(key_resolver, movedRanges, src, dest); - - ResolutionSplitRequest req; - req.front = range.second; - req.offset = amount; - req.range = range.first; - - ResolutionSplitReply split = - wait(brokenPromiseToNever(self->resolvers[metrics.lastItem()->second].split.getReply( - req, TaskPriority::ResolutionMetrics))); - KeyRangeRef moveRange = range.second ? KeyRangeRef(range.first.begin, split.key) - : KeyRangeRef(split.key, range.first.end); - movedRanges.push_back_deep(movedRanges.arena(), ResolverMoveRef(moveRange, dest)); - TraceEvent("MovingResolutionRange") - .detail("Src", src) - .detail("Dest", dest) - .detail("Amount", amount) - .detail("StartRange", range.first) - .detail("MoveRange", moveRange) - .detail("Used", split.used) - .detail("KeyResolverRanges", key_resolver.size()); - amount -= split.used; - if (moveRange != range.first || amount <= 0) - break; - } - for (auto& it : movedRanges) - key_resolver.insert(it.range, it.dest); - // for(auto& it : key_resolver.ranges()) - // TraceEvent("KeyResolver").detail("Range", it.range()).detail("Value", it.value()); - - self->resolverChangesVersion = self->version + 1; - for (auto& p : self->commitProxies) - self->resolverNeedingChanges.insert(p.id()); - self->resolverChanges.set(movedRanges); - } catch (Error& e) { - if (e.code() != error_code_operation_failed) - throw; - } - } - } -} - ACTOR Future changeCoordinators(Reference self) { loop { ChangeCoordinatorsRequest req = waitNext(self->clusterController.changeCoordinators.getFuture()); @@ -1127,8 +988,8 @@ ACTOR Future>> recruitEverything( newTLogServers(self, recruits, oldLogSystem, &confChanges)); // Update recovery related information to the newly elected sequencer (master) process. - wait(brokenPromiseToNever(self->masterInterface.updateRecoveryData.getReply( - UpdateRecoveryDataRequest(self->recoveryTransactionVersion, self->lastEpochEnd, self->commitProxies)))); + wait(brokenPromiseToNever(self->masterInterface.updateRecoveryData.getReply(UpdateRecoveryDataRequest( + self->recoveryTransactionVersion, self->lastEpochEnd, self->commitProxies, self->resolvers)))); return confChanges; } @@ -1802,14 +1663,6 @@ ACTOR Future clusterRecoveryCore(Reference self) { .detail("RecoveryDuration", recoveryDuration) .trackLatest(self->clusterRecoveryStateEventHolder->trackingKey); - TraceEvent(getRecoveryEventName(ClusterRecoveryEventType::CLUSTER_RECOVERY_AVAILABLE_EVENT_NAME).c_str(), - self->dbgid) - .detail("AvailableAtVersion", self->version) - .trackLatest(self->clusterRecoveryAvailableEventHolder->trackingKey); - - if (self->resolvers.size() > 1) - self->addActor.send(resolutionBalancing(self)); - self->addActor.send(changeCoordinators(self)); Database cx = openDBOnServer(self->dbInfo, TaskPriority::DefaultEndpoint, LockAware::True); self->addActor.send(configurationMonitor(self, cx)); diff --git a/fdbserver/ClusterRecovery.actor.h b/fdbserver/ClusterRecovery.actor.h index 36dcb1bed9..19911fdd02 100644 --- a/fdbserver/ClusterRecovery.actor.h +++ b/fdbserver/ClusterRecovery.actor.h @@ -185,7 +185,6 @@ struct ClusterRecoveryData : NonCopyable, ReferenceCounted ServerCoordinators coordinators; Reference logSystem; - Version version; // The last version assigned to a proxy by getVersion() double lastVersionTime; LogSystemDiskQueueAdapter* txnStateLogAdapter; IKeyValueStore* txnStateStore; @@ -225,10 +224,6 @@ struct ClusterRecoveryData : NonCopyable, ReferenceCounted RecoveryState recoveryState; - AsyncVar>> resolverChanges; - Version resolverChangesVersion; - std::set resolverNeedingChanges; - PromiseStream> addActor; Reference> recruitmentStalled; bool forceRecovery; @@ -266,12 +261,11 @@ struct ClusterRecoveryData : NonCopyable, ReferenceCounted : controllerData(controllerData), dbgid(masterInterface.id()), lastEpochEnd(invalidVersion), recoveryTransactionVersion(invalidVersion), lastCommitTime(0), liveCommittedVersion(invalidVersion), databaseLocked(false), minKnownCommittedVersion(invalidVersion), hasConfiguration(false), - coordinators(coordinators), version(invalidVersion), lastVersionTime(0), txnStateStore(nullptr), - memoryLimit(2e9), dbId(dbId), masterInterface(masterInterface), masterLifetime(masterLifetimeToken), - clusterController(clusterController), cstate(coordinators, addActor, dbgid), dbInfo(dbInfo), - registrationCount(0), addActor(addActor), recruitmentStalled(makeReference>(false)), - forceRecovery(forceRecovery), neverCreated(false), safeLocality(tagLocalityInvalid), - primaryLocality(tagLocalityInvalid), cc("Master", dbgid.toString()), + coordinators(coordinators), lastVersionTime(0), txnStateStore(nullptr), memoryLimit(2e9), dbId(dbId), + masterInterface(masterInterface), masterLifetime(masterLifetimeToken), clusterController(clusterController), + cstate(coordinators, addActor, dbgid), dbInfo(dbInfo), registrationCount(0), addActor(addActor), + recruitmentStalled(makeReference>(false)), forceRecovery(forceRecovery), neverCreated(false), + safeLocality(tagLocalityInvalid), primaryLocality(tagLocalityInvalid), cc("Master", dbgid.toString()), changeCoordinatorsRequests("ChangeCoordinatorsRequests", cc), getCommitVersionRequests("GetCommitVersionRequests", cc), backupWorkerDoneRequests("BackupWorkerDoneRequests", cc), diff --git a/fdbserver/MasterInterface.h b/fdbserver/MasterInterface.h index 90d49e9492..1b7918a583 100644 --- a/fdbserver/MasterInterface.h +++ b/fdbserver/MasterInterface.h @@ -23,14 +23,15 @@ #pragma once #include "fdbclient/CommitProxyInterface.h" -#include "fdbclient/FDBTypes.h" -#include "fdbclient/StorageServerInterface.h" #include "fdbclient/CommitTransaction.h" #include "fdbclient/DatabaseConfiguration.h" -#include "fdbserver/TLogInterface.h" +#include "fdbclient/FDBTypes.h" #include "fdbclient/Notified.h" +#include "fdbclient/StorageServerInterface.h" +#include "fdbserver/ResolverInterface.h" +#include "fdbserver/TLogInterface.h" -typedef uint64_t DBRecoveryCount; +using DBRecoveryCount = uint64_t; struct MasterInterface { constexpr static FileIdentifier file_identifier = 5979145; @@ -155,18 +156,20 @@ struct UpdateRecoveryDataRequest { Version recoveryTransactionVersion; Version lastEpochEnd; std::vector commitProxies; + std::vector resolvers; ReplyPromise reply; - UpdateRecoveryDataRequest() {} + UpdateRecoveryDataRequest() = default; UpdateRecoveryDataRequest(Version recoveryTransactionVersion, Version lastEpochEnd, - std::vector commitProxies) + const std::vector& commitProxies, + const std::vector& resolvers) : recoveryTransactionVersion(recoveryTransactionVersion), lastEpochEnd(lastEpochEnd), - commitProxies(commitProxies) {} + commitProxies(commitProxies), resolvers(resolvers) {} template void serialize(Ar& ar) { - serializer(ar, recoveryTransactionVersion, lastEpochEnd, commitProxies, reply); + serializer(ar, recoveryTransactionVersion, lastEpochEnd, commitProxies, resolvers, reply); } }; diff --git a/fdbserver/masterserver.actor.cpp b/fdbserver/masterserver.actor.cpp index d9e860eb52..5a43adf927 100644 --- a/fdbserver/masterserver.actor.cpp +++ b/fdbserver/masterserver.actor.cpp @@ -67,6 +67,9 @@ struct MasterData : NonCopyable, ReferenceCounted { std::vector commitProxies; std::map lastCommitProxyVersionReplies; + std::vector resolvers; + + PromiseStream> addActor; MasterInterface myInterface; @@ -94,7 +97,7 @@ struct MasterData : NonCopyable, ReferenceCounted { : dbgid(myInterface.id()), lastEpochEnd(invalidVersion), recoveryTransactionVersion(invalidVersion), liveCommittedVersion(invalidVersion), databaseLocked(false), minKnownCommittedVersion(invalidVersion), coordinators(coordinators), version(invalidVersion), lastVersionTime(0), txnStateStore(nullptr), - myInterface(myInterface), forceRecovery(forceRecovery), cc("Master", dbgid.toString()), + addActor(addActor), myInterface(myInterface), forceRecovery(forceRecovery), cc("Master", dbgid.toString()), getCommitVersionRequests("GetCommitVersionRequests", cc), getLiveCommittedVersionRequests("GetLiveCommittedVersionRequests", cc), reportLiveCommittedVersionRequests("ReportLiveCommittedVersionRequests", cc) { @@ -110,6 +113,145 @@ struct MasterData : NonCopyable, ReferenceCounted { } }; +static std::pair findRange(CoalescedKeyRangeMap& key_resolver, + Standalone>& movedRanges, + int src, + int dest) { + auto ranges = key_resolver.ranges(); + auto prev = ranges.begin(); + auto it = ranges.begin(); + ++it; + if (it == ranges.end()) { + if (ranges.begin().value() != src || + std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(ranges.begin()->range(), dest)) != + movedRanges.end()) + throw operation_failed(); + return std::make_pair(ranges.begin().range(), true); + } + + std::set borders; + // If possible expand an existing boundary between the two resolvers + for (; it != ranges.end(); ++it) { + if (it->value() == src && prev->value() == dest && + std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(it->range(), dest)) == + movedRanges.end()) { + return std::make_pair(it->range(), true); + } + if (it->value() == dest && prev->value() == src && + std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(prev->range(), dest)) == + movedRanges.end()) { + return std::make_pair(prev->range(), false); + } + if (it->value() == dest) + borders.insert(prev->value()); + if (prev->value() == dest) + borders.insert(it->value()); + ++prev; + } + + prev = ranges.begin(); + it = ranges.begin(); + ++it; + // If possible create a new boundry which doesn't exist yet + for (; it != ranges.end(); ++it) { + if (it->value() == src && !borders.count(prev->value()) && + std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(it->range(), dest)) == + movedRanges.end()) { + return std::make_pair(it->range(), true); + } + if (prev->value() == src && !borders.count(it->value()) && + std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(prev->range(), dest)) == + movedRanges.end()) { + return std::make_pair(prev->range(), false); + } + ++prev; + } + + it = ranges.begin(); + for (; it != ranges.end(); ++it) { + if (it->value() == src && + std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(it->range(), dest)) == + movedRanges.end()) { + return std::make_pair(it->range(), true); + } + } + throw operation_failed(); // we are already attempting to move all of the data one resolver is assigned, so do not + // move anything +} + +// Balance key ranges among resolvers so that their load are evenly distributed. +ACTOR Future resolutionBalancing(Reference self) { + state CoalescedKeyRangeMap key_resolver; + key_resolver.insert(allKeys, 0); + loop { + wait(delay(SERVER_KNOBS->MIN_BALANCE_TIME, TaskPriority::ResolutionMetrics)); + while (self->resolverChanges.get().size()) + wait(self->resolverChanges.onChange()); + state std::vector> futures; + for (auto& p : self->resolvers) + futures.push_back( + brokenPromiseToNever(p.metrics.getReply(ResolutionMetricsRequest(), TaskPriority::ResolutionMetrics))); + wait(waitForAll(futures)); + state IndexedSet, NoMetric> metrics; + + int64_t total = 0; + for (int i = 0; i < futures.size(); i++) { + total += futures[i].get().value; + metrics.insert(std::make_pair(futures[i].get().value, i), NoMetric()); + //TraceEvent("ResolverMetric").detail("I", i).detail("Metric", futures[i].get()); + } + if (metrics.lastItem()->first - metrics.begin()->first > SERVER_KNOBS->MIN_BALANCE_DIFFERENCE) { + try { + state int src = metrics.lastItem()->second; + state int dest = metrics.begin()->second; + state int64_t amount = std::min(metrics.lastItem()->first - total / self->resolvers.size(), + total / self->resolvers.size() - metrics.begin()->first) / + 2; + state Standalone> movedRanges; + + loop { + state std::pair range = findRange(key_resolver, movedRanges, src, dest); + + ResolutionSplitRequest req; + req.front = range.second; + req.offset = amount; + req.range = range.first; + + ResolutionSplitReply split = + wait(brokenPromiseToNever(self->resolvers[metrics.lastItem()->second].split.getReply( + req, TaskPriority::ResolutionMetrics))); + KeyRangeRef moveRange = range.second ? KeyRangeRef(range.first.begin, split.key) + : KeyRangeRef(split.key, range.first.end); + movedRanges.push_back_deep(movedRanges.arena(), ResolverMoveRef(moveRange, dest)); + TraceEvent("MovingResolutionRange") + .detail("Src", src) + .detail("Dest", dest) + .detail("Amount", amount) + .detail("StartRange", range.first) + .detail("MoveRange", moveRange) + .detail("Used", split.used) + .detail("KeyResolverRanges", key_resolver.size()); + amount -= split.used; + if (moveRange != range.first || amount <= 0) + break; + } + for (auto& it : movedRanges) + key_resolver.insert(it.range, it.dest); + // for(auto& it : key_resolver.ranges()) + // TraceEvent("KeyResolver").detail("Range", it.range()).detail("Value", it.value()); + + self->resolverChangesVersion = self->version + 1; + for (auto& p : self->commitProxies) + self->resolverNeedingChanges.insert(p.id()); + self->resolverChanges.set(movedRanges); + } catch (Error& e) { + if (e.code() != error_code_operation_failed) + throw; + } + } + } +} + ACTOR Future getVersion(Reference self, GetCommitVersionRequest req) { state Span span("M:getVersion"_loc, { req.spanContext }); state std::map::iterator proxyItr = @@ -244,31 +386,33 @@ ACTOR Future serveLiveCommittedVersion(Reference self) { ACTOR Future updateRecoveryData(Reference self) { loop { - choose { - when(UpdateRecoveryDataRequest req = waitNext(self->myInterface.updateRecoveryData.getFuture())) { - TraceEvent("UpdateRecoveryData", self->dbgid) - .detail("RecoveryTxnVersion", req.recoveryTransactionVersion) - .detail("LastEpochEnd", req.lastEpochEnd) - .detail("NumCommitProxies", req.commitProxies.size()); + UpdateRecoveryDataRequest req = waitNext(self->myInterface.updateRecoveryData.getFuture()); + TraceEvent("UpdateRecoveryData", self->dbgid) + .detail("RecoveryTxnVersion", req.recoveryTransactionVersion) + .detail("LastEpochEnd", req.lastEpochEnd) + .detail("NumCommitProxies", req.commitProxies.size()); - if (self->recoveryTransactionVersion == invalidVersion || - req.recoveryTransactionVersion > self->recoveryTransactionVersion) { - self->recoveryTransactionVersion = req.recoveryTransactionVersion; - } - if (self->lastEpochEnd == invalidVersion || req.lastEpochEnd > self->lastEpochEnd) { - self->lastEpochEnd = req.lastEpochEnd; - } - if (req.commitProxies.size() > 0) { - self->commitProxies = req.commitProxies; - self->lastCommitProxyVersionReplies.clear(); + if (self->recoveryTransactionVersion == invalidVersion || + req.recoveryTransactionVersion > self->recoveryTransactionVersion) { + self->recoveryTransactionVersion = req.recoveryTransactionVersion; + } + if (self->lastEpochEnd == invalidVersion || req.lastEpochEnd > self->lastEpochEnd) { + self->lastEpochEnd = req.lastEpochEnd; + } + if (req.commitProxies.size() > 0) { + self->commitProxies = req.commitProxies; + self->lastCommitProxyVersionReplies.clear(); - for (auto& p : self->commitProxies) { - self->lastCommitProxyVersionReplies[p.id()] = CommitProxyVersionReplies(); - } - } - req.reply.send(Void()); + for (auto& p : self->commitProxies) { + self->lastCommitProxyVersionReplies[p.id()] = CommitProxyVersionReplies(); } } + + self->resolvers = req.resolvers; + if (req.resolvers.size() > 1) + self->addActor.send(resolutionBalancing(self)); + + req.reply.send(Void()); } } diff --git a/fdbserver/workloads/KillRegion.actor.cpp b/fdbserver/workloads/KillRegion.actor.cpp index d15255bd7b..08b3f62946 100644 --- a/fdbserver/workloads/KillRegion.actor.cpp +++ b/fdbserver/workloads/KillRegion.actor.cpp @@ -107,7 +107,10 @@ struct KillRegionWorkload : TestWorkload { DatabaseConfiguration conf = wait(getDatabaseConfiguration(cx)); - TraceEvent("ForceRecovery_GotConfig").detail("Conf", conf.toString()); + TraceEvent("ForceRecovery_GotConfig") + .setMaxEventLength(11000) + .setMaxFieldLength(10000) + .detail("Conf", conf.toString()); if (conf.usableRegions > 1) { loop { From 213e37191c7baaae9c45fbd5f8298b450340f0c3 Mon Sep 17 00:00:00 2001 From: Jingyu Zhou Date: Fri, 18 Mar 2022 15:27:31 -0700 Subject: [PATCH 2/6] Add code coverage macro for resolution balancing --- fdbserver/masterserver.actor.cpp | 1 + 1 file changed, 1 insertion(+) diff --git a/fdbserver/masterserver.actor.cpp b/fdbserver/masterserver.actor.cpp index 5a43adf927..f9898073c2 100644 --- a/fdbserver/masterserver.actor.cpp +++ b/fdbserver/masterserver.actor.cpp @@ -306,6 +306,7 @@ ACTOR Future getVersion(Reference self, GetCommitVersionReques rep.resolverChangesVersion = self->resolverChangesVersion; self->resolverNeedingChanges.erase(req.requestingProxy); + TEST(!rep.resolverChanges.empty()); // resolution balancing moves keyranges if (self->resolverNeedingChanges.empty()) self->resolverChanges.set(Standalone>()); } From 7736ea87b05047f6d64b8eb6a9be97e0700482cd Mon Sep 17 00:00:00 2001 From: Jingyu Zhou Date: Fri, 18 Mar 2022 15:57:34 -0700 Subject: [PATCH 3/6] Fix code format --- fdbserver/masterserver.actor.cpp | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/fdbserver/masterserver.actor.cpp b/fdbserver/masterserver.actor.cpp index f9898073c2..33dacd45ad 100644 --- a/fdbserver/masterserver.actor.cpp +++ b/fdbserver/masterserver.actor.cpp @@ -389,12 +389,12 @@ ACTOR Future updateRecoveryData(Reference self) { loop { UpdateRecoveryDataRequest req = waitNext(self->myInterface.updateRecoveryData.getFuture()); TraceEvent("UpdateRecoveryData", self->dbgid) - .detail("RecoveryTxnVersion", req.recoveryTransactionVersion) - .detail("LastEpochEnd", req.lastEpochEnd) - .detail("NumCommitProxies", req.commitProxies.size()); + .detail("RecoveryTxnVersion", req.recoveryTransactionVersion) + .detail("LastEpochEnd", req.lastEpochEnd) + .detail("NumCommitProxies", req.commitProxies.size()); if (self->recoveryTransactionVersion == invalidVersion || - req.recoveryTransactionVersion > self->recoveryTransactionVersion) { + req.recoveryTransactionVersion > self->recoveryTransactionVersion) { self->recoveryTransactionVersion = req.recoveryTransactionVersion; } if (self->lastEpochEnd == invalidVersion || req.lastEpochEnd > self->lastEpochEnd) { From 6e8d16538dcb4b0dd0e9c5f3ec71f64306b4bde4 Mon Sep 17 00:00:00 2001 From: Jingyu Zhou Date: Mon, 21 Mar 2022 13:48:09 -0700 Subject: [PATCH 4/6] Remove an unused variable from MasterData --- fdbserver/masterserver.actor.cpp | 10 +++------- 1 file changed, 3 insertions(+), 7 deletions(-) diff --git a/fdbserver/masterserver.actor.cpp b/fdbserver/masterserver.actor.cpp index 33dacd45ad..f79680c888 100644 --- a/fdbserver/masterserver.actor.cpp +++ b/fdbserver/masterserver.actor.cpp @@ -63,7 +63,6 @@ struct MasterData : NonCopyable, ReferenceCounted { Version version; // The last version assigned to a proxy by getVersion() double lastVersionTime; - IKeyValueStore* txnStateStore; std::vector commitProxies; std::map lastCommitProxyVersionReplies; @@ -96,8 +95,8 @@ struct MasterData : NonCopyable, ReferenceCounted { : dbgid(myInterface.id()), lastEpochEnd(invalidVersion), recoveryTransactionVersion(invalidVersion), liveCommittedVersion(invalidVersion), databaseLocked(false), minKnownCommittedVersion(invalidVersion), - coordinators(coordinators), version(invalidVersion), lastVersionTime(0), txnStateStore(nullptr), - addActor(addActor), myInterface(myInterface), forceRecovery(forceRecovery), cc("Master", dbgid.toString()), + coordinators(coordinators), version(invalidVersion), lastVersionTime(0), addActor(addActor), + myInterface(myInterface), forceRecovery(forceRecovery), cc("Master", dbgid.toString()), getCommitVersionRequests("GetCommitVersionRequests", cc), getLiveCommittedVersionRequests("GetLiveCommittedVersionRequests", cc), reportLiveCommittedVersionRequests("ReportLiveCommittedVersionRequests", cc) { @@ -107,10 +106,7 @@ struct MasterData : NonCopyable, ReferenceCounted { forceRecovery = false; } } - ~MasterData() { - if (txnStateStore) - txnStateStore->close(); - } + ~MasterData() = default; }; static std::pair findRange(CoalescedKeyRangeMap& key_resolver, From 437e7d27c643407ba57365eca38f050c8e557da1 Mon Sep 17 00:00:00 2001 From: Jingyu Zhou Date: Mon, 21 Mar 2022 16:38:23 -0700 Subject: [PATCH 5/6] Use trigger to start resolutionBalancing Trigger when there are more than one resolvers. This also avoids the problem of receiving multiple UpdateRecoveryDataRequests. --- fdbserver/masterserver.actor.cpp | 35 +++++++++----------------------- 1 file changed, 10 insertions(+), 25 deletions(-) diff --git a/fdbserver/masterserver.actor.cpp b/fdbserver/masterserver.actor.cpp index f79680c888..4e7c0b1653 100644 --- a/fdbserver/masterserver.actor.cpp +++ b/fdbserver/masterserver.actor.cpp @@ -20,31 +20,15 @@ #include -#include "fdbclient/NativeAPI.actor.h" -#include "fdbclient/Notified.h" -#include "fdbclient/SystemData.h" -#include "fdbrpc/FailureMonitor.h" -#include "fdbrpc/PerfMetric.h" #include "fdbrpc/sim_validation.h" -#include "fdbrpc/simulator.h" -#include "fdbserver/ApplyMetadataMutation.h" -#include "fdbserver/BackupProgress.actor.h" -#include "fdbserver/ConflictSet.h" #include "fdbserver/CoordinatedState.h" #include "fdbserver/CoordinationInterface.h" // copy constructors for ServerCoordinators class -#include "fdbserver/DBCoreState.h" -#include "fdbserver/DataDistribution.actor.h" -#include "fdbserver/IKeyValueStore.h" #include "fdbserver/Knobs.h" -#include "fdbserver/LogSystemDiskQueueAdapter.h" #include "fdbserver/MasterInterface.h" -#include "fdbserver/ProxyCommitData.actor.h" -#include "fdbserver/RecoveryState.h" #include "fdbserver/ServerDBInfo.h" -#include "fdbserver/WaitFailure.h" -#include "fdbserver/WorkerInterface.actor.h" #include "flow/ActorCollection.h" #include "flow/Trace.h" +#include "flow/genericactors.actor.h" #include "flow/actorcompiler.h" // This must be the last #include. @@ -68,13 +52,12 @@ struct MasterData : NonCopyable, ReferenceCounted { std::map lastCommitProxyVersionReplies; std::vector resolvers; - PromiseStream> addActor; - MasterInterface myInterface; AsyncVar>> resolverChanges; Version resolverChangesVersion; std::set resolverNeedingChanges; + AsyncTrigger triggerResolution; bool forceRecovery; @@ -90,13 +73,12 @@ struct MasterData : NonCopyable, ReferenceCounted { ServerCoordinators const& coordinators, ClusterControllerFullInterface const& clusterController, Standalone const& dbId, - PromiseStream> const& addActor, bool forceRecovery) : dbgid(myInterface.id()), lastEpochEnd(invalidVersion), recoveryTransactionVersion(invalidVersion), liveCommittedVersion(invalidVersion), databaseLocked(false), minKnownCommittedVersion(invalidVersion), - coordinators(coordinators), version(invalidVersion), lastVersionTime(0), addActor(addActor), - myInterface(myInterface), forceRecovery(forceRecovery), cc("Master", dbgid.toString()), + coordinators(coordinators), version(invalidVersion), lastVersionTime(0), myInterface(myInterface), + forceRecovery(forceRecovery), cc("Master", dbgid.toString()), getCommitVersionRequests("GetCommitVersionRequests", cc), getLiveCommittedVersionRequests("GetLiveCommittedVersionRequests", cc), reportLiveCommittedVersionRequests("ReportLiveCommittedVersionRequests", cc) { @@ -177,6 +159,8 @@ static std::pair findRange(CoalescedKeyRangeMap& key_res // Balance key ranges among resolvers so that their load are evenly distributed. ACTOR Future resolutionBalancing(Reference self) { + wait(self->triggerResolution.onTrigger()); + state CoalescedKeyRangeMap key_resolver; key_resolver.insert(allKeys, 0); loop { @@ -407,7 +391,7 @@ ACTOR Future updateRecoveryData(Reference self) { self->resolvers = req.resolvers; if (req.resolvers.size() > 1) - self->addActor.send(resolutionBalancing(self)); + self->triggerResolution.trigger(); req.reply.send(Void()); } @@ -454,14 +438,15 @@ ACTOR Future masterServer(MasterInterface mi, state Future onDBChange = Void(); state PromiseStream> addActor; - state Reference self(new MasterData( - db, mi, coordinators, db->get().clusterInterface, LiteralStringRef(""), addActor, forceRecovery)); + state Reference self( + new MasterData(db, mi, coordinators, db->get().clusterInterface, LiteralStringRef(""), forceRecovery)); state Future collection = actorCollection(addActor.getFuture()); addActor.send(traceRole(Role::MASTER, mi.id())); addActor.send(provideVersions(self)); addActor.send(serveLiveCommittedVersion(self)); addActor.send(updateRecoveryData(self)); + addActor.send(resolutionBalancing(self)); TEST(!lifetime.isStillValid(db->get().masterLifetime, mi.id() == db->get().master.id())); // Master born doomed TraceEvent("MasterLifetime", self->dbgid).detail("LifetimeToken", lifetime.toString()); From 0c88be03931656d1c9e25c35413891cab0bea876 Mon Sep 17 00:00:00 2001 From: Jingyu Zhou Date: Mon, 21 Mar 2022 21:35:48 -0700 Subject: [PATCH 6/6] Refactor resolution balancing into separate files --- fdbserver/CMakeLists.txt | 2 + fdbserver/ResolutionBalancer.actor.cpp | 187 +++++++++++++++++++++++++ fdbserver/ResolutionBalancer.actor.h | 64 +++++++++ fdbserver/masterserver.actor.cpp | 186 ++---------------------- 4 files changed, 266 insertions(+), 173 deletions(-) create mode 100644 fdbserver/ResolutionBalancer.actor.cpp create mode 100644 fdbserver/ResolutionBalancer.actor.h diff --git a/fdbserver/CMakeLists.txt b/fdbserver/CMakeLists.txt index d16d65b51d..98d9d3eeae 100644 --- a/fdbserver/CMakeLists.txt +++ b/fdbserver/CMakeLists.txt @@ -96,6 +96,8 @@ set(FDBSERVER_SRCS Ratekeeper.h RatekeeperInterface.h RecoveryState.h + ResolutionBalancer.actor.cpp + ResolutionBalancer.actor.h Resolver.actor.cpp ResolverInterface.h RestoreApplier.actor.cpp diff --git a/fdbserver/ResolutionBalancer.actor.cpp b/fdbserver/ResolutionBalancer.actor.cpp new file mode 100644 index 0000000000..6e569d7058 --- /dev/null +++ b/fdbserver/ResolutionBalancer.actor.cpp @@ -0,0 +1,187 @@ +/* + * ResolutionBalancer.actor.cpp + * + * This source file is part of the FoundationDB open source project + * + * Copyright 2013-2022 Apple Inc. and the FoundationDB project authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +#include "fdbserver/ResolutionBalancer.actor.h" + +#include "fdbclient/KeyRangeMap.h" +#include "fdbserver/MasterInterface.h" +#include "fdbserver/Knobs.h" +#include "flow/flow.h" + +#include "flow/actorcompiler.h" // This must be the last #include. + +void ResolutionBalancer::setResolvers(const std::vector& v) { + resolvers = v; + if (resolvers.size() > 1) + triggerResolution.trigger(); +} + +void ResolutionBalancer::setChangesInReply(UID requestingProxy, GetCommitVersionReply& rep) { + if (resolverNeedingChanges.count(requestingProxy)) { + rep.resolverChanges = resolverChanges.get(); + rep.resolverChangesVersion = resolverChangesVersion; + resolverNeedingChanges.erase(requestingProxy); + + TEST(!rep.resolverChanges.empty()); // resolution balancing moves keyranges + if (resolverNeedingChanges.empty()) + resolverChanges.set(Standalone>()); + } +} + +static std::pair findRange(CoalescedKeyRangeMap& key_resolver, + Standalone>& movedRanges, + int src, + int dest) { + auto ranges = key_resolver.ranges(); + auto prev = ranges.begin(); + auto it = ranges.begin(); + ++it; + if (it == ranges.end()) { + if (ranges.begin().value() != src || + std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(ranges.begin()->range(), dest)) != + movedRanges.end()) + throw operation_failed(); + return std::make_pair(ranges.begin().range(), true); + } + + std::set borders; + // If possible expand an existing boundary between the two resolvers + for (; it != ranges.end(); ++it) { + if (it->value() == src && prev->value() == dest && + std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(it->range(), dest)) == + movedRanges.end()) { + return std::make_pair(it->range(), true); + } + if (it->value() == dest && prev->value() == src && + std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(prev->range(), dest)) == + movedRanges.end()) { + return std::make_pair(prev->range(), false); + } + if (it->value() == dest) + borders.insert(prev->value()); + if (prev->value() == dest) + borders.insert(it->value()); + ++prev; + } + + prev = ranges.begin(); + it = ranges.begin(); + ++it; + // If possible create a new boundry which doesn't exist yet + for (; it != ranges.end(); ++it) { + if (it->value() == src && !borders.count(prev->value()) && + std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(it->range(), dest)) == + movedRanges.end()) { + return std::make_pair(it->range(), true); + } + if (prev->value() == src && !borders.count(it->value()) && + std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(prev->range(), dest)) == + movedRanges.end()) { + return std::make_pair(prev->range(), false); + } + ++prev; + } + + it = ranges.begin(); + for (; it != ranges.end(); ++it) { + if (it->value() == src && + std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(it->range(), dest)) == + movedRanges.end()) { + return std::make_pair(it->range(), true); + } + } + throw operation_failed(); // we are already attempting to move all of the data one resolver is assigned, so do not + // move anything +} + +// Balance key ranges among resolvers so that their load are evenly distributed. +ACTOR Future ResolutionBalancer::resolutionBalancing_impl(ResolutionBalancer* self) { + wait(self->triggerResolution.onTrigger()); + + state CoalescedKeyRangeMap key_resolver; + key_resolver.insert(allKeys, 0); + loop { + wait(delay(SERVER_KNOBS->MIN_BALANCE_TIME, TaskPriority::ResolutionMetrics)); + while (self->resolverChanges.get().size()) + wait(self->resolverChanges.onChange()); + state std::vector> futures; + for (auto& p : self->resolvers) + futures.push_back( + brokenPromiseToNever(p.metrics.getReply(ResolutionMetricsRequest(), TaskPriority::ResolutionMetrics))); + wait(waitForAll(futures)); + state IndexedSet, NoMetric> metrics; + + int64_t total = 0; + for (int i = 0; i < futures.size(); i++) { + total += futures[i].get().value; + metrics.insert(std::make_pair(futures[i].get().value, i), NoMetric()); + //TraceEvent("ResolverMetric").detail("I", i).detail("Metric", futures[i].get()); + } + if (metrics.lastItem()->first - metrics.begin()->first > SERVER_KNOBS->MIN_BALANCE_DIFFERENCE) { + try { + state int src = metrics.lastItem()->second; + state int dest = metrics.begin()->second; + state int64_t amount = std::min(metrics.lastItem()->first - total / self->resolvers.size(), + total / self->resolvers.size() - metrics.begin()->first) / + 2; + state Standalone> movedRanges; + + loop { + state std::pair range = findRange(key_resolver, movedRanges, src, dest); + + ResolutionSplitRequest req; + req.front = range.second; + req.offset = amount; + req.range = range.first; + + ResolutionSplitReply split = + wait(brokenPromiseToNever(self->resolvers[metrics.lastItem()->second].split.getReply( + req, TaskPriority::ResolutionMetrics))); + KeyRangeRef moveRange = range.second ? KeyRangeRef(range.first.begin, split.key) + : KeyRangeRef(split.key, range.first.end); + movedRanges.push_back_deep(movedRanges.arena(), ResolverMoveRef(moveRange, dest)); + TraceEvent("MovingResolutionRange") + .detail("Src", src) + .detail("Dest", dest) + .detail("Amount", amount) + .detail("StartRange", range.first) + .detail("MoveRange", moveRange) + .detail("Used", split.used) + .detail("KeyResolverRanges", key_resolver.size()); + amount -= split.used; + if (moveRange != range.first || amount <= 0) + break; + } + for (auto& it : movedRanges) + key_resolver.insert(it.range, it.dest); + // for(auto& it : key_resolver.ranges()) + // TraceEvent("KeyResolver").detail("Range", it.range()).detail("Value", it.value()); + + self->resolverChangesVersion = *self->pVersion + 1; + for (auto& p : self->commitProxies) + self->resolverNeedingChanges.insert(p.id()); + self->resolverChanges.set(movedRanges); + } catch (Error& e) { + if (e.code() != error_code_operation_failed) + throw; + } + } + } +} diff --git a/fdbserver/ResolutionBalancer.actor.h b/fdbserver/ResolutionBalancer.actor.h new file mode 100644 index 0000000000..263195f0e5 --- /dev/null +++ b/fdbserver/ResolutionBalancer.actor.h @@ -0,0 +1,64 @@ +/* + * ResolutionBalancer.actor.h + * + * This source file is part of the FoundationDB open source project + * + * Copyright 2013-2022 Apple Inc. and the FoundationDB project authors + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +#include "fdbclient/CommitProxyInterface.h" +#include "fdbserver/ResolverInterface.h" +#if defined(NO_INTELLISENSE) && !defined(FDBSERVER_RESOLUTION_BALANCER_G_H) +#define FDBSERVER_RESOLUTION_BALANCER_G_H +#include "fdbserver/ResolutionBalancer.actor.g.h" +#elif !defined(FDBSERVER_RESOLUTION_BALANCER_H) +#define FDBSERVER_RESOLUTION_BALANCER_H + +#include + +#include "fdbclient/FDBTypes.h" +#include "fdbserver/MasterInterface.h" +#include "flow/Arena.h" +#include "flow/IRandom.h" +#include "flow/genericactors.actor.h" + +struct ResolutionBalancer { + AsyncVar>> resolverChanges; + Version resolverChangesVersion = invalidVersion; + std::set resolverNeedingChanges; + + Version* pVersion; // points to MasterData::version + + std::vector commitProxies; + std::vector resolvers; + AsyncTrigger triggerResolution; + + ResolutionBalancer(Version* version) : pVersion(version) {} + + Future resolutionBalancing() { return resolutionBalancing_impl(this); } + + ACTOR static Future resolutionBalancing_impl(ResolutionBalancer* self); + + // Sets resolver interfaces. Trigger resolutionBalancing() actor if more + // than one resolvers are present. + void setResolvers(const std::vector& resolvers); + + void setCommitProxies(const std::vector& proxies) { commitProxies = proxies; } + + void setChangesInReply(UID requestingProxy, GetCommitVersionReply& rep); +}; + +#include "flow/unactorcompiler.h" +#endif diff --git a/fdbserver/masterserver.actor.cpp b/fdbserver/masterserver.actor.cpp index 4e7c0b1653..7a2ba7b2de 100644 --- a/fdbserver/masterserver.actor.cpp +++ b/fdbserver/masterserver.actor.cpp @@ -25,10 +25,10 @@ #include "fdbserver/CoordinationInterface.h" // copy constructors for ServerCoordinators class #include "fdbserver/Knobs.h" #include "fdbserver/MasterInterface.h" +#include "fdbserver/ResolutionBalancer.actor.h" #include "fdbserver/ServerDBInfo.h" #include "flow/ActorCollection.h" #include "flow/Trace.h" -#include "flow/genericactors.actor.h" #include "flow/actorcompiler.h" // This must be the last #include. @@ -48,16 +48,11 @@ struct MasterData : NonCopyable, ReferenceCounted { Version version; // The last version assigned to a proxy by getVersion() double lastVersionTime; - std::vector commitProxies; std::map lastCommitProxyVersionReplies; - std::vector resolvers; MasterInterface myInterface; - AsyncVar>> resolverChanges; - Version resolverChangesVersion; - std::set resolverNeedingChanges; - AsyncTrigger triggerResolution; + ResolutionBalancer resolutionBalancer; bool forceRecovery; @@ -67,6 +62,7 @@ struct MasterData : NonCopyable, ReferenceCounted { Counter reportLiveCommittedVersionRequests; Future logger; + Future balancer; MasterData(Reference const> const& dbInfo, MasterInterface const& myInterface, @@ -78,7 +74,7 @@ struct MasterData : NonCopyable, ReferenceCounted { : dbgid(myInterface.id()), lastEpochEnd(invalidVersion), recoveryTransactionVersion(invalidVersion), liveCommittedVersion(invalidVersion), databaseLocked(false), minKnownCommittedVersion(invalidVersion), coordinators(coordinators), version(invalidVersion), lastVersionTime(0), myInterface(myInterface), - forceRecovery(forceRecovery), cc("Master", dbgid.toString()), + resolutionBalancer(&version), forceRecovery(forceRecovery), cc("Master", dbgid.toString()), getCommitVersionRequests("GetCommitVersionRequests", cc), getLiveCommittedVersionRequests("GetLiveCommittedVersionRequests", cc), reportLiveCommittedVersionRequests("ReportLiveCommittedVersionRequests", cc) { @@ -87,151 +83,11 @@ struct MasterData : NonCopyable, ReferenceCounted { TraceEvent(SevError, "ForcedRecoveryRequiresDcID").log(); forceRecovery = false; } + balancer = resolutionBalancer.resolutionBalancing(); } ~MasterData() = default; }; -static std::pair findRange(CoalescedKeyRangeMap& key_resolver, - Standalone>& movedRanges, - int src, - int dest) { - auto ranges = key_resolver.ranges(); - auto prev = ranges.begin(); - auto it = ranges.begin(); - ++it; - if (it == ranges.end()) { - if (ranges.begin().value() != src || - std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(ranges.begin()->range(), dest)) != - movedRanges.end()) - throw operation_failed(); - return std::make_pair(ranges.begin().range(), true); - } - - std::set borders; - // If possible expand an existing boundary between the two resolvers - for (; it != ranges.end(); ++it) { - if (it->value() == src && prev->value() == dest && - std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(it->range(), dest)) == - movedRanges.end()) { - return std::make_pair(it->range(), true); - } - if (it->value() == dest && prev->value() == src && - std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(prev->range(), dest)) == - movedRanges.end()) { - return std::make_pair(prev->range(), false); - } - if (it->value() == dest) - borders.insert(prev->value()); - if (prev->value() == dest) - borders.insert(it->value()); - ++prev; - } - - prev = ranges.begin(); - it = ranges.begin(); - ++it; - // If possible create a new boundry which doesn't exist yet - for (; it != ranges.end(); ++it) { - if (it->value() == src && !borders.count(prev->value()) && - std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(it->range(), dest)) == - movedRanges.end()) { - return std::make_pair(it->range(), true); - } - if (prev->value() == src && !borders.count(it->value()) && - std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(prev->range(), dest)) == - movedRanges.end()) { - return std::make_pair(prev->range(), false); - } - ++prev; - } - - it = ranges.begin(); - for (; it != ranges.end(); ++it) { - if (it->value() == src && - std::find(movedRanges.begin(), movedRanges.end(), ResolverMoveRef(it->range(), dest)) == - movedRanges.end()) { - return std::make_pair(it->range(), true); - } - } - throw operation_failed(); // we are already attempting to move all of the data one resolver is assigned, so do not - // move anything -} - -// Balance key ranges among resolvers so that their load are evenly distributed. -ACTOR Future resolutionBalancing(Reference self) { - wait(self->triggerResolution.onTrigger()); - - state CoalescedKeyRangeMap key_resolver; - key_resolver.insert(allKeys, 0); - loop { - wait(delay(SERVER_KNOBS->MIN_BALANCE_TIME, TaskPriority::ResolutionMetrics)); - while (self->resolverChanges.get().size()) - wait(self->resolverChanges.onChange()); - state std::vector> futures; - for (auto& p : self->resolvers) - futures.push_back( - brokenPromiseToNever(p.metrics.getReply(ResolutionMetricsRequest(), TaskPriority::ResolutionMetrics))); - wait(waitForAll(futures)); - state IndexedSet, NoMetric> metrics; - - int64_t total = 0; - for (int i = 0; i < futures.size(); i++) { - total += futures[i].get().value; - metrics.insert(std::make_pair(futures[i].get().value, i), NoMetric()); - //TraceEvent("ResolverMetric").detail("I", i).detail("Metric", futures[i].get()); - } - if (metrics.lastItem()->first - metrics.begin()->first > SERVER_KNOBS->MIN_BALANCE_DIFFERENCE) { - try { - state int src = metrics.lastItem()->second; - state int dest = metrics.begin()->second; - state int64_t amount = std::min(metrics.lastItem()->first - total / self->resolvers.size(), - total / self->resolvers.size() - metrics.begin()->first) / - 2; - state Standalone> movedRanges; - - loop { - state std::pair range = findRange(key_resolver, movedRanges, src, dest); - - ResolutionSplitRequest req; - req.front = range.second; - req.offset = amount; - req.range = range.first; - - ResolutionSplitReply split = - wait(brokenPromiseToNever(self->resolvers[metrics.lastItem()->second].split.getReply( - req, TaskPriority::ResolutionMetrics))); - KeyRangeRef moveRange = range.second ? KeyRangeRef(range.first.begin, split.key) - : KeyRangeRef(split.key, range.first.end); - movedRanges.push_back_deep(movedRanges.arena(), ResolverMoveRef(moveRange, dest)); - TraceEvent("MovingResolutionRange") - .detail("Src", src) - .detail("Dest", dest) - .detail("Amount", amount) - .detail("StartRange", range.first) - .detail("MoveRange", moveRange) - .detail("Used", split.used) - .detail("KeyResolverRanges", key_resolver.size()); - amount -= split.used; - if (moveRange != range.first || amount <= 0) - break; - } - for (auto& it : movedRanges) - key_resolver.insert(it.range, it.dest); - // for(auto& it : key_resolver.ranges()) - // TraceEvent("KeyResolver").detail("Range", it.range()).detail("Value", it.value()); - - self->resolverChangesVersion = self->version + 1; - for (auto& p : self->commitProxies) - self->resolverNeedingChanges.insert(p.id()); - self->resolverChanges.set(movedRanges); - } catch (Error& e) { - if (e.code() != error_code_operation_failed) - throw; - } - } - } -} - ACTOR Future getVersion(Reference self, GetCommitVersionRequest req) { state Span span("M:getVersion"_loc, { req.spanContext }); state std::map::iterator proxyItr = @@ -281,15 +137,7 @@ ACTOR Future getVersion(Reference self, GetCommitVersionReques TEST(maxVersionGap); // Maximum possible version gap self->lastVersionTime = t1; - if (self->resolverNeedingChanges.count(req.requestingProxy)) { - rep.resolverChanges = self->resolverChanges.get(); - rep.resolverChangesVersion = self->resolverChangesVersion; - self->resolverNeedingChanges.erase(req.requestingProxy); - - TEST(!rep.resolverChanges.empty()); // resolution balancing moves keyranges - if (self->resolverNeedingChanges.empty()) - self->resolverChanges.set(Standalone>()); - } + self->resolutionBalancer.setChangesInReply(req.requestingProxy, rep); } rep.version = self->version; @@ -311,16 +159,11 @@ ACTOR Future getVersion(Reference self, GetCommitVersionReques ACTOR Future provideVersions(Reference self) { state ActorCollection versionActors(false); - for (auto& p : self->commitProxies) - self->lastCommitProxyVersionReplies[p.id()] = CommitProxyVersionReplies(); - - loop { - choose { - when(GetCommitVersionRequest req = waitNext(self->myInterface.getCommitVersion.getFuture())) { - versionActors.add(getVersion(self, req)); - } - when(wait(versionActors.getResult())) {} + loop choose { + when(GetCommitVersionRequest req = waitNext(self->myInterface.getCommitVersion.getFuture())) { + versionActors.add(getVersion(self, req)); } + when(wait(versionActors.getResult())) {} } } @@ -381,17 +224,15 @@ ACTOR Future updateRecoveryData(Reference self) { self->lastEpochEnd = req.lastEpochEnd; } if (req.commitProxies.size() > 0) { - self->commitProxies = req.commitProxies; self->lastCommitProxyVersionReplies.clear(); - for (auto& p : self->commitProxies) { + for (auto& p : req.commitProxies) { self->lastCommitProxyVersionReplies[p.id()] = CommitProxyVersionReplies(); } } - self->resolvers = req.resolvers; - if (req.resolvers.size() > 1) - self->triggerResolution.trigger(); + self->resolutionBalancer.setCommitProxies(req.commitProxies); + self->resolutionBalancer.setResolvers(req.resolvers); req.reply.send(Void()); } @@ -446,7 +287,6 @@ ACTOR Future masterServer(MasterInterface mi, addActor.send(provideVersions(self)); addActor.send(serveLiveCommittedVersion(self)); addActor.send(updateRecoveryData(self)); - addActor.send(resolutionBalancing(self)); TEST(!lifetime.isStillValid(db->get().masterLifetime, mi.id() == db->get().master.id())); // Master born doomed TraceEvent("MasterLifetime", self->dbgid).detail("LifetimeToken", lifetime.toString());