diff --git a/fdbserver/datadistributor/DDTeamCollection.actor.cpp b/fdbserver/datadistributor/DDTeamCollection.actor.cpp index a5fefc814c..3f5d6c4cfc 100644 --- a/fdbserver/datadistributor/DDTeamCollection.actor.cpp +++ b/fdbserver/datadistributor/DDTeamCollection.actor.cpp @@ -19,6 +19,7 @@ */ #include +#include #include "fdbclient/SystemData.h" #include "fdbrpc/simulator.h" @@ -30,6 +31,7 @@ #include "TCInfo.h" #include "ExclusionTracker.h" #include "flow/IRandom.h" +#include "flow/ScopeExit.h" #include "flow/Trace.h" #include "flow/network.h" #include "flow/TxnCounters.h" @@ -1035,9 +1037,11 @@ public: bool lastHealthy{ false }; bool lastOptimal{ false }; bool lastWrongConfiguration = team->isWrongConfiguration(); + bool lastContainsFailed = false; bool trackHealthyTeam = team->size() == self->configuration.storageTeamSize; bool lastZeroHealthy = self->zeroHealthyTeams->get(); bool firstCheck = true; + std::unordered_set submittedShards; Future zeroServerLeftLogger; @@ -1105,9 +1109,9 @@ public: team->setHealthy(healthy); // Unhealthy teams won't be chosen by bestTeam bool optimal = team->isOptimal() && healthy; bool containsFailed = self->teamContainsFailedServer(team); - bool retryUnhealthyShards = serversLeft < team->size() && - self->shardsAffectedByTeamFailure->hasShards( - ShardsAffectedByTeamFailure::Team(team->getServerIDs(), self->primary)); + bool retryUnhealthyShards = + !healthy && self->shardsAffectedByTeamFailure->hasShards( + ShardsAffectedByTeamFailure::Team(team->getServerIDs(), self->primary)); if (retryUnhealthyShards) { // Partial moves can leave a merged shard associated with this team without another health change. change.push_back(delay(SERVER_KNOBS->CHECK_TEAM_DELAY, TaskPriority::DataDistributionLow)); @@ -1148,9 +1152,11 @@ public: lastOptimal = optimal; } - if (serversLeft != lastServersLeft || anyUndesired != lastAnyUndesired || - anyWrongConfiguration != lastWrongConfiguration || anyWigglingServer != lastAnyWigglingServer || - recheck) { // NOTE: do not check wrongSize + bool teamStateChanged = serversLeft != lastServersLeft || anyUndesired != lastAnyUndesired || + anyWrongConfiguration != lastWrongConfiguration || + anyWigglingServer != lastAnyWigglingServer || + containsFailed != lastContainsFailed; + if (teamStateChanged || recheck) { // NOTE: do not check wrongSize if (logTeamEvents) { TraceEvent("ServerTeamHealthChanged", self->distributorId) .suppressFor(1.0) @@ -1204,6 +1210,7 @@ public: lastAnyUndesired = anyUndesired; lastWrongConfiguration = anyWrongConfiguration; lastAnyWigglingServer = anyWigglingServer; + lastContainsFailed = containsFailed; int lastPriority = team->getPriority(); if (team->size() == 0) { @@ -1265,6 +1272,20 @@ public: std::vector shards = self->shardsAffectedByTeamFailure->getShardsFor( ShardsAffectedByTeamFailure::Team(team->getServerIDs(), self->primary)); + if (teamStateChanged || !retryUnhealthyShards) { + submittedShards.clear(); + } else { + // An unchanged range may still be waiting behind the relocation pipeline gate. Only retry + // newly mapped ranges until the team state changes, and forget ranges that disappeared. + std::unordered_set mappedShards(shards.begin(), shards.end()); + for (auto it = submittedShards.begin(); it != submittedShards.end();) { + if (!mappedShards.contains(*it)) { + it = submittedShards.erase(it); + } else { + ++it; + } + } + } TraceEvent(SevVerbose, "ServerTeamRelocatingShards", self->distributorId) .detail("Info", team->getDesc()) @@ -1272,6 +1293,9 @@ public: .detail("Shards", shards.size()); for (int i = 0; i < shards.size(); i++) { + if (retryUnhealthyShards && !submittedShards.insert(shards[i]).second) { + continue; + } // Make it high priority to move keys off failed server or else RelocateShards may never be // addressed int maxPriority = containsFailed ? SERVER_KNOBS->PRIORITY_TEAM_FAILED : team->getPriority(); @@ -7132,6 +7156,86 @@ public: recruitment.cancel(); co_await delay(0); } + + static Future TeamTracker_RetriesMergedShardForUndesiredServer() { + auto serverKnobs = const_cast(SERVER_KNOBS); + const double originalCheckTeamDelay = serverKnobs->CHECK_TEAM_DELAY; + serverKnobs->CHECK_TEAM_DELAY = 0.05; + auto restoreCheckTeamDelay = ScopeExit( + [serverKnobs, originalCheckTeamDelay]() { serverKnobs->CHECK_TEAM_DELAY = originalCheckTeamDelay; }); + + auto shards = makeReference(); + shards->setCheckMode(ShardsAffectedByTeamFailure::CheckMode::ForceCheck); + + const UID undesired(1, 0), leftServer(2, 0), rightServer(3, 0), healthy1(4, 0), healthy2(5, 0); + const ShardsAffectedByTeamFailure::Team left({ undesired, leftServer }, true); + const ShardsAffectedByTeamFailure::Team right({ undesired, rightServer }, true); + const ShardsAffectedByTeamFailure::Team healthy({ healthy1, healthy2 }, true); + const KeyRange mergedRange = KeyRangeRef("a"_sr, "c"_sr); + const KeyRange leftRange = KeyRangeRef("a"_sr, "b"_sr); + const KeyRange rightRange = KeyRangeRef("b"_sr, "c"_sr); + + shards->assignRangeToTeams(leftRange, { left }); + shards->assignRangeToTeams(rightRange, { right }); + + Reference policy = makeReference(2, "zoneid", makeReference()); + auto collection = testTeamCollection(2, policy, 5, shards); + collection->teamCollections = { collection.get() }; + collection->initialFailureReactionDelay = Future(Void()); + collection->server_status.set( + undesired, + ServerStatus(IsFailed::False, + IsUndesired::True, + IsWiggling::False, + collection->server_info[undesired]->getLastKnownInterface().locality)); + FutureStream relocations = collection->output.getFuture(); + + collection->addTeam(std::set({ undesired, leftServer }), IsInitialTeam::True); + collection->addTeam(std::set({ undesired, rightServer }), IsInitialTeam::True); + collection->addTeam(std::set({ healthy1, healthy2 }), IsInitialTeam::True); + co_await delay(0.1); + + ASSERT(relocations.isReady()); + RelocateShard initialLeft = relocations.pop(); + ASSERT(relocations.isReady()); + RelocateShard initialRight = relocations.pop(); + ASSERT(!relocations.isReady()); + ASSERT((initialLeft.keys == leftRange && initialRight.keys == rightRange) || + (initialLeft.keys == rightRange && initialRight.keys == leftRange)); + + shards->defineShard(mergedRange); + shards->moveShard(leftRange, { healthy }); + shards->finishMove(leftRange); + shards->moveShard(rightRange, { healthy }); + shards->finishMove(rightRange); + ASSERT_EQ(shards->getNumberOfShards(undesired), 2); + + co_await delay(SERVER_KNOBS->CHECK_TEAM_DELAY + 0.1); + ASSERT(relocations.isReady()); + RelocateShard retryLeft = relocations.pop(); + ASSERT(relocations.isReady()); + RelocateShard retryRight = relocations.pop(); + ASSERT(!relocations.isReady()); + ASSERT(retryLeft.keys == mergedRange); + ASSERT(retryRight.keys == mergedRange); + + co_await delay(SERVER_KNOBS->CHECK_TEAM_DELAY + 0.1); + ASSERT(!relocations.isReady()); + + NetworkAddress failedAddress = collection->server_info[undesired]->getLastKnownInterface().address(); + collection->excludedServers.set(AddressExclusion(failedAddress.ip, failedAddress.port), + DDTeamCollection::Status::FAILED); + co_await delay(SERVER_KNOBS->CHECK_TEAM_DELAY + 0.1); + ASSERT(relocations.isReady()); + RelocateShard failedLeft = relocations.pop(); + ASSERT(relocations.isReady()); + RelocateShard failedRight = relocations.pop(); + ASSERT(!relocations.isReady()); + ASSERT(failedLeft.keys == mergedRange); + ASSERT(failedRight.keys == mergedRange); + ASSERT_EQ(failedLeft.priority, SERVER_KNOBS->PRIORITY_TEAM_FAILED); + ASSERT_EQ(failedRight.priority, SERVER_KNOBS->PRIORITY_TEAM_FAILED); + } }; TEST_CASE("DataDistribution/AddTeamsBestOf/UseMachineID") { @@ -7295,3 +7399,8 @@ TEST_CASE("/DataDistribution/Recruitment/RecruitmentFailedCooldownReleasesId") { wait(DDTeamCollectionUnitTest::InitializeStorage_RecruitmentFailedCooldownReleasesId()); return Void(); } + +TEST_CASE("/DataDistribution/TeamTracker/RetriesMergedShardForUndesiredServer") { + wait(DDTeamCollectionUnitTest::TeamTracker_RetriesMergedShardForUndesiredServer()); + return Void(); +}