diff options
Diffstat (limited to 'src/mongo/db/s/balancer/balancer_policy.cpp')
| -rw-r--r-- | src/mongo/db/s/balancer/balancer_policy.cpp | 450 |
1 files changed, 152 insertions, 298 deletions
diff --git a/src/mongo/db/s/balancer/balancer_policy.cpp b/src/mongo/db/s/balancer/balancer_policy.cpp index 4e11ab57cb8..ec4041b9667 100644 --- a/src/mongo/db/s/balancer/balancer_policy.cpp +++ b/src/mongo/db/s/balancer/balancer_policy.cpp @@ -58,134 +58,65 @@ namespace { // optimal average across all shards for a zone for a rebalancing migration to be initiated. const size_t kDefaultImbalanceThreshold = 1; -ChunkType makeChunkType(const NamespaceString& nss, const Chunk& chunk) { - ChunkType ct{nss, chunk.getRange(), chunk.getLastmod(), chunk.getShardId()}; - ct.setJumbo(chunk.isJumbo()); - return ct; -} - -/** - * Return a vector of zones after they have been normalized according to the given chunk - * configuration. - * - * If a zone covers only partially a chunk, boundaries of that zone will be shrank so that the - * normalized zone won't overlap with that chunk. The boundaries of a normalized zone will never - * fall in the middle of a chunk. - * - * Additionally the vector will contain also zones for the "NoZone", - */ -std::vector<ZoneRange> normalizeZones(const ChunkManager& cm, const ZoneInfo& zoneInfo) { - std::vector<ZoneRange> normalizedRanges; - - auto lastMax = cm.getShardKeyPattern().getKeyPattern().globalMin(); - - for (const auto& [max, zoneRange] : zoneInfo.zoneRanges()) { - const auto& minChunk = cm.findIntersectingChunkWithSimpleCollation(zoneRange.min); - const auto gtMin = - SimpleBSONObjComparator::kInstance.evaluate(zoneRange.min > minChunk.getMin()); - const auto& normalizedMin = gtMin ? minChunk.getMax() : zoneRange.min; - - - const auto& maxChunk = cm.findIntersectingChunkWithSimpleCollation(zoneRange.max); - const auto gtMax = - SimpleBSONObjComparator::kInstance.evaluate(zoneRange.max > maxChunk.getMin()) && - SimpleBSONObjComparator::kInstance.evaluate( - zoneRange.max != cm.getShardKeyPattern().getKeyPattern().globalMax()); - const auto& normalizedMax = gtMax ? maxChunk.getMin() : zoneRange.max; - - - if (SimpleBSONObjComparator::kInstance.evaluate(normalizedMin == normalizedMax)) { - // This zone does not fully contain any chunk thus we can ignore it - continue; - } +} // namespace - if (SimpleBSONObjComparator::kInstance.evaluate(normalizedMin != lastMax)) { - // The zone is not contiguous with the previous one so we add a kNoZoneRange - // does not fully contain any chunk so we will ignore it - normalizedRanges.emplace_back(lastMax, normalizedMin, ZoneInfo::kNoZoneName); - } +DistributionStatus::DistributionStatus(NamespaceString nss, + ShardToChunksMap shardToChunksMap, + ZoneInfo zoneInfo) + : _nss(std::move(nss)), + _shardChunks(std::move(shardToChunksMap)), + _zoneInfo(std::move(zoneInfo)) {} - normalizedRanges.emplace_back(normalizedMin, normalizedMax, zoneRange.zone); - lastMax = normalizedMax; - } +size_t DistributionStatus::totalChunks() const { + size_t total = 0; - const auto& globalMaxKey = cm.getShardKeyPattern().getKeyPattern().globalMax(); - if (SimpleBSONObjComparator::kInstance.evaluate(lastMax != globalMaxKey)) { - normalizedRanges.emplace_back(lastMax, globalMaxKey, ZoneInfo::kNoZoneName); + for (const auto& shardChunk : _shardChunks) { + total += shardChunk.second.size(); } - return normalizedRanges; -} - -} // namespace -DistributionStatus::DistributionStatus(NamespaceString nss, - ZoneInfo zoneInfo, - const ChunkManager& chunkMngr) - : _nss(std::move(nss)), _zoneInfo(std::move(zoneInfo)), _chunkMngr(chunkMngr) { - - _normalizedZones = normalizeZones(_chunkMngr, _zoneInfo); - - for (const auto& zoneRange : _normalizedZones) { - chunkMngr.forEachOverlappingChunk( - zoneRange.min, zoneRange.max, false /* isMaxInclusive */, [&](const auto& chunkInfo) { - _shardToZoneSizeMap[chunkInfo.getShardId()][zoneRange.zone]++; - return true; - }); - } + return total; } size_t DistributionStatus::totalChunksWithTag(const std::string& tag) const { size_t total = 0; - for (const auto& [_, zoneSizeMap] : _shardToZoneSizeMap) { - const auto& zoneIt = zoneSizeMap.find(tag); - if (zoneIt != zoneSizeMap.end()) { - total += zoneIt->second; - } + + for (const auto& shardChunk : _shardChunks) { + total += numberOfChunksInShardWithTag(shardChunk.first, tag); } + return total; } size_t DistributionStatus::numberOfChunksInShard(const ShardId& shardId) const { - const auto shardZonesIt = _shardToZoneSizeMap.find(shardId); - if (shardZonesIt == _shardToZoneSizeMap.end()) { - return 0; - } - size_t total = 0; - for (const auto& [_, numChunks] : shardZonesIt->second) { - total += numChunks; - } - return total; + const auto& shardChunks = getChunks(shardId); + return shardChunks.size(); } size_t DistributionStatus::numberOfChunksInShardWithTag(const ShardId& shardId, const string& tag) const { - const auto shardZonesIt = _shardToZoneSizeMap.find(shardId); - if (shardZonesIt == _shardToZoneSizeMap.end()) { - return 0; - } - const auto& shardTags = shardZonesIt->second; + const auto& shardChunks = getChunks(shardId); - const auto& zoneIt = shardTags.find(tag); - if (zoneIt == shardTags.end()) { - return 0; + size_t total = 0; + + for (const auto& chunk : shardChunks) { + if (tag == getTagForChunk(chunk)) { + total++; + } } - return zoneIt->second; -} -string DistributionStatus::getTagForRange(const ChunkRange& range) const { - return _zoneInfo.getZoneForChunk(range); + return total; } -const StringMap<size_t>& DistributionStatus::getChunksPerTagMap(const ShardId& shardId) const { - static const StringMap<size_t> emptyMap; - const auto shardZonesIt = _shardToZoneSizeMap.find(shardId); - if (shardZonesIt == _shardToZoneSizeMap.end()) { - return emptyMap; - } - return shardZonesIt->second; +const vector<ChunkType>& DistributionStatus::getChunks(const ShardId& shardId) const { + ShardToChunksMap::const_iterator i = _shardChunks.find(shardId); + invariant(i != _shardChunks.end()); + + return i->second; } -const string ZoneInfo::kNoZoneName = ""; +string DistributionStatus::getTagForChunk(const ChunkType& chunk) const { + return _zoneInfo.getZoneForChunk(chunk.getRange()); +} ZoneInfo::ZoneInfo() : _zoneRanges(SimpleBSONObjComparator::kInstance.makeBSONObjIndexedMap<ZoneRange>()) {} @@ -237,11 +168,11 @@ string ZoneInfo::getZoneForChunk(const ChunkRange& chunk) const { // We should never have a partial overlap with a chunk range. If it happens, treat it as if this // chunk doesn't belong to a tag if (minIntersect != maxIntersect) { - return ZoneInfo::kNoZoneName; + return ""; } if (minIntersect == _zoneRanges.end()) { - return ZoneInfo::kNoZoneName; + return ""; } const ZoneRange& intersectRange = minIntersect->second; @@ -252,7 +183,7 @@ string ZoneInfo::getZoneForChunk(const ChunkRange& chunk) const { return intersectRange.zone; } - return ZoneInfo::kNoZoneName; + return ""; } void DistributionStatus::report(BSONObjBuilder* builder) const { @@ -260,15 +191,15 @@ void DistributionStatus::report(BSONObjBuilder* builder) const { // Report all shards BSONArrayBuilder shardArr(builder->subarrayStart("shards")); - for (const auto& [shardId, zoneSizeMap] : _shardToZoneSizeMap) { + for (const auto& shardChunk : _shardChunks) { BSONObjBuilder shardEntry(shardArr.subobjStart()); - shardEntry.append("name", shardId.toString()); + shardEntry.append("name", shardChunk.first.toString()); - BSONObjBuilder tagsObj(shardEntry.subobjStart("tags")); - for (const auto& [tagName, numChunks] : zoneSizeMap) { - tagsObj.appendNumber(tagName, static_cast<long long>(numChunks)); + BSONArrayBuilder chunkArr(shardEntry.subarrayStart("chunks")); + for (const auto& chunk : shardChunk.second) { + chunkArr.append(chunk.toConfigBSON()); } - tagsObj.doneFast(); + chunkArr.doneFast(); shardEntry.doneFast(); } @@ -311,7 +242,7 @@ Status BalancerPolicy::isShardSuitableReceiver(const ClusterStatistics::ShardSta str::stream() << stat.shardId << " is currently draining."}; } - if (chunkTag != ZoneInfo::kNoZoneName && !stat.shardTags.count(chunkTag)) { + if (!chunkTag.empty() && !stat.shardTags.count(chunkTag)) { return {ErrorCodes::IllegalOperation, str::stream() << stat.shardId << " is not in the correct zone " << chunkTag}; } @@ -429,28 +360,10 @@ MigrateInfo chooseRandomMigration(const ShardStatisticsVector& shardStats, "fromShardId"_attr = sourceShardId, "toShardId"_attr = destShardId); - const auto& randomChunk = [&] { - const auto numChunksOnSourceShard = distribution.numberOfChunksInShard(sourceShardId); - const auto rndChunkIdx = getRandomIndex(numChunksOnSourceShard); - ChunkType rndChunk; - - int idx{0}; - distribution.getChunkManager().forEachChunk([&](const auto& chunk) { - if (chunk.getShardId() == sourceShardId && idx++ == rndChunkIdx) { - rndChunk = makeChunkType(distribution.nss(), chunk); - rndChunk.setJumbo(chunk.isJumbo()); - return false; - } - return true; - }); - - invariant(rndChunk.getShard().isValid()); - return rndChunk; - }(); - + const auto& chunks = distribution.getChunks(sourceShardId); return {destShardId, - randomChunk, + chunks[getRandomIndex(chunks.size())], MoveChunkRequest::ForceJumbo::kDoNotForce, MigrateInfo::chunksImbalance}; } @@ -482,70 +395,57 @@ vector<MigrateInfo> BalancerPolicy::balance(const ShardStatisticsVector& shardSt if (!availableShards->count(stat.shardId)) continue; + const vector<ChunkType>& chunks = distribution.getChunks(stat.shardId); + + if (chunks.empty()) + continue; + // Now we know we need to move to chunks off this shard, but only if permitted by the // tags policy unsigned numJumboChunks = 0; - const auto& chunksPerTagMap = distribution.getChunksPerTagMap(stat.shardId); - for (const auto& tagIt : chunksPerTagMap) { - const auto& zoneName = tagIt.first; - for (const auto& zoneRange : distribution.getNormalizedZones()) { - if (zoneRange.zone != zoneName) { - continue; - } + // Since we have to move all chunks, lets just do in order + for (const auto& chunk : chunks) { + if (chunk.getJumbo()) { + numJumboChunks++; + continue; + } - distribution.getChunkManager().forEachOverlappingChunk( - zoneRange.min, - zoneRange.max, - false /* isMaxInclusive */, - [&](const auto& chunk) { - if (chunk.getShardId() != stat.shardId) { - return true; // continue - } - if (chunk.isJumbo()) { - numJumboChunks++; - return true; // continue - } - - const ShardId to = _getLeastLoadedReceiverShard( - shardStats, distribution, zoneName, *availableShards); - if (!to.isValid()) { - if (migrations.empty()) { - LOGV2_WARNING( - 21889, - "Chunk {chunk} is on a draining shard, but no appropriate " - "recipient found", - "Chunk is on a draining shard, but no appropriate " - "recipient found", - "chunk"_attr = redact( - makeChunkType(distribution.nss(), chunk).toString())); - } - return true; // continue - } - invariant(to != stat.shardId); - - migrations.emplace_back(to, - makeChunkType(distribution.nss(), chunk), - MoveChunkRequest::ForceJumbo::kForceBalancer, - MigrateInfo::drain); - invariant(availableShards->erase(stat.shardId)); - invariant(availableShards->erase(to)); - return false; // break - }); + const string tag = distribution.getTagForChunk(chunk); + const ShardId to = + _getLeastLoadedReceiverShard(shardStats, distribution, tag, *availableShards); + if (!to.isValid()) { if (migrations.empty()) { - LOGV2_WARNING(21890, - "Unable to find any chunk to move from draining shard " - "{shardId}. numJumboChunks: {numJumboChunks}", - "Unable to find any chunk to move from draining shard", - "shardId"_attr = stat.shardId, - "numJumboChunks"_attr = numJumboChunks); - } - - if (availableShards->size() < 2) { - return migrations; + LOGV2_WARNING(21889, + "Chunk {chunk} is on a draining shard, but no appropriate " + "recipient found", + "Chunk is on a draining shard, but no appropriate " + "recipient found", + "chunk"_attr = redact(chunk.toString())); } + continue; } + + invariant(to != stat.shardId); + migrations.emplace_back( + to, chunk, MoveChunkRequest::ForceJumbo::kForceBalancer, MigrateInfo::drain); + invariant(availableShards->erase(stat.shardId)); + invariant(availableShards->erase(to)); + break; + } + + if (migrations.empty()) { + LOGV2_WARNING(21890, + "Unable to find any chunk to move from draining shard " + "{shardId}. numJumboChunks: {numJumboChunks}", + "Unable to find any chunk to move from draining shard", + "shardId"_attr = stat.shardId, + "numJumboChunks"_attr = numJumboChunks); + } + + if (availableShards->size() < 2) { + return migrations; } } } @@ -553,77 +453,57 @@ vector<MigrateInfo> BalancerPolicy::balance(const ShardStatisticsVector& shardSt // 2) Check for chunks, which are on the wrong shard and must be moved off of it if (!distribution.tags().empty()) { for (const auto& stat : shardStats) { - if (!availableShards->count(stat.shardId)) continue; - const auto& chunksPerTagMap = distribution.getChunksPerTagMap(stat.shardId); - for (const auto& tagIt : chunksPerTagMap) { - const auto& zoneName = tagIt.first; + const vector<ChunkType>& chunks = distribution.getChunks(stat.shardId); - if (zoneName == ZoneInfo::kNoZoneName) + for (const auto& chunk : chunks) { + const string tag = distribution.getTagForChunk(chunk); + + if (tag.empty()) continue; - if (stat.shardTags.count(zoneName)) + if (stat.shardTags.count(tag)) continue; - for (const auto& zoneRange : distribution.getNormalizedZones()) { - if (zoneRange.zone != zoneName) { - continue; - } - distribution.getChunkManager().forEachOverlappingChunk( - zoneRange.min, - zoneRange.max, - false /* isMaxInclusive */, - [&](const auto& chunk) { - if (chunk.getShardId() != stat.shardId) { - return true; // continue - } - if (chunk.isJumbo()) { - LOGV2_WARNING( - 21891, - "Chunk {chunk} violates zone {zone}, but it is jumbo and " - "cannot be " - "moved", - "Chunk violates zone, but it is jumbo and cannot be moved", - "chunk"_attr = - redact(makeChunkType(distribution.nss(), chunk).toString()), - "zone"_attr = redact(zoneName)); - return true; // continue - } - - const ShardId to = _getLeastLoadedReceiverShard( - shardStats, distribution, zoneName, *availableShards); - if (!to.isValid()) { - if (migrations.empty()) { - LOGV2_WARNING( - 21892, - "Chunk {chunk} violates zone {zone}, but no appropriate " - "recipient found", - "Chunk violates zone, but no appropriate recipient found", - "chunk"_attr = redact( - makeChunkType(distribution.nss(), chunk).toString()), - "zone"_attr = redact(zoneName)); - } - return true; // continue - } - invariant(to != stat.shardId); - - migrations.emplace_back( - to, - makeChunkType(distribution.nss(), chunk), - forceJumbo ? MoveChunkRequest::ForceJumbo::kForceBalancer - : MoveChunkRequest::ForceJumbo::kDoNotForce, - MigrateInfo::zoneViolation); - invariant(availableShards->erase(stat.shardId)); - invariant(availableShards->erase(to)); - return false; // break - }); + if (chunk.getJumbo()) { + LOGV2_WARNING( + 21891, + "Chunk {chunk} violates zone {zone}, but it is jumbo and cannot be moved", + "Chunk violates zone, but it is jumbo and cannot be moved", + "chunk"_attr = redact(chunk.toString()), + "zone"_attr = redact(tag)); + continue; } - if (availableShards->size() < 2) { - return migrations; + const ShardId to = + _getLeastLoadedReceiverShard(shardStats, distribution, tag, *availableShards); + if (!to.isValid()) { + if (migrations.empty()) { + LOGV2_WARNING(21892, + "Chunk {chunk} violates zone {zone}, but no appropriate " + "recipient found", + "Chunk violates zone, but no appropriate recipient found", + "chunk"_attr = redact(chunk.toString()), + "zone"_attr = redact(tag)); + } + continue; } + + invariant(to != stat.shardId); + migrations.emplace_back(to, + chunk, + forceJumbo ? MoveChunkRequest::ForceJumbo::kForceBalancer + : MoveChunkRequest::ForceJumbo::kDoNotForce, + MigrateInfo::zoneViolation); + invariant(availableShards->erase(stat.shardId)); + invariant(availableShards->erase(to)); + break; + } + + if (availableShards->size() < 2) { + return migrations; } } } @@ -631,35 +511,24 @@ vector<MigrateInfo> BalancerPolicy::balance(const ShardStatisticsVector& shardSt // 3) for each tag balance vector<string> tagsPlusEmpty(distribution.tags().begin(), distribution.tags().end()); - tagsPlusEmpty.push_back(ZoneInfo::kNoZoneName); + tagsPlusEmpty.push_back(""); for (const auto& tag : tagsPlusEmpty) { + const size_t totalNumberOfChunksWithTag = + (tag.empty() ? distribution.totalChunks() : distribution.totalChunksWithTag(tag)); - const auto totalNumberOfChunksWithTag = [&] { - if (tag == ZoneInfo::kNoZoneName) { - return static_cast<size_t>(distribution.getChunkManager().numChunks()); - } - return distribution.totalChunksWithTag(tag); - }(); + size_t totalNumberOfShardsWithTag = 0; - const auto totalNumberOfShardsWithTag = [&] { - if (tag == ZoneInfo::kNoZoneName) { - return shardStats.size(); - } - - size_t numShardsWithTag{0}; - for (const auto& stat : shardStats) { - if (stat.shardTags.count(tag)) { - numShardsWithTag++; - } + for (const auto& stat : shardStats) { + if (tag.empty() || stat.shardTags.count(tag)) { + totalNumberOfShardsWithTag++; } - return numShardsWithTag; - }(); + } // Skip zones which have no shards assigned to them. This situation is not harmful, but // should not be possible so warn the operator to correct it. if (totalNumberOfShardsWithTag == 0) { - if (tag != ZoneInfo::kNoZoneName) { + if (!tag.empty()) { LOGV2_WARNING( 21893, "Zone {zone} in collection {namespace} has no assigned shards and chunks " @@ -696,7 +565,7 @@ boost::optional<MigrateInfo> BalancerPolicy::balanceSingleChunk( const ChunkType& chunk, const ShardStatisticsVector& shardStats, const DistributionStatus& distribution) { - const string tag = distribution.getTagForRange(chunk.getRange()); + const string tag = distribution.getTagForChunk(chunk); stdx::unordered_set<ShardId> availableShards; std::transform(shardStats.begin(), @@ -773,41 +642,26 @@ bool BalancerPolicy::_singleZoneBalance(const ShardStatisticsVector& shardStats, if (imbalance < kDefaultImbalanceThreshold) return false; + const vector<ChunkType>& chunks = distribution.getChunks(from); + unsigned numJumboChunks = 0; - bool chunkFound = false; - for (const auto& zoneRange : distribution.getNormalizedZones()) { - if (zoneRange.zone != tag) { + for (const auto& chunk : chunks) { + if (distribution.getTagForChunk(chunk) != tag) continue; - } - - distribution.getChunkManager().forEachOverlappingChunk( - zoneRange.min, zoneRange.max, false /* isMaxInclusive */, [&](const auto& chunk) { - if (chunk.getShardId() != from) { - return true; // continue - } - if (chunk.isJumbo()) { - numJumboChunks++; - return true; // continue - } - - migrations->emplace_back(to, - makeChunkType(distribution.nss(), chunk), - forceJumbo, - MigrateInfo::chunksImbalance); - invariant(availableShards->erase(chunk.getShardId())); - invariant(availableShards->erase(to)); - chunkFound = true; - return false; // break - }); - - if (chunkFound) { - return chunkFound; + if (chunk.getJumbo()) { + numJumboChunks++; + continue; } + + migrations->emplace_back(to, chunk, forceJumbo, MigrateInfo::chunksImbalance); + invariant(availableShards->erase(chunk.getShard())); + invariant(availableShards->erase(to)); + return true; } - if (!chunkFound && numJumboChunks) { + if (numJumboChunks) { LOGV2_WARNING( 21894, "Shard: {shardId}, collection: {namespace} has only jumbo chunks for " @@ -819,7 +673,7 @@ bool BalancerPolicy::_singleZoneBalance(const ShardStatisticsVector& shardStats, "numJumboChunks"_attr = numJumboChunks); } - return chunkFound; + return false; } ZoneRange::ZoneRange(const BSONObj& a_min, const BSONObj& a_max, const std::string& _zone) |
