diff options
Diffstat (limited to 'src/mongo/db/s/config/initial_split_policy.cpp')
| -rw-r--r-- | src/mongo/db/s/config/initial_split_policy.cpp | 108 |
1 files changed, 39 insertions, 69 deletions
diff --git a/src/mongo/db/s/config/initial_split_policy.cpp b/src/mongo/db/s/config/initial_split_policy.cpp index 15c5a345c59..5bdfb55d6c3 100644 --- a/src/mongo/db/s/config/initial_split_policy.cpp +++ b/src/mongo/db/s/config/initial_split_policy.cpp @@ -37,14 +37,12 @@ #include "mongo/db/bson/dotted_path_support.h" #include "mongo/db/catalog/collection_catalog.h" #include "mongo/db/curop.h" -#include "mongo/db/pipeline/document_source.h" #include "mongo/db/pipeline/lite_parsed_pipeline.h" #include "mongo/db/pipeline/process_interface/shardsvr_process_interface.h" #include "mongo/db/pipeline/sharded_agg_helpers.h" #include "mongo/db/s/balancer/balancer_policy.h" #include "mongo/db/s/sharding_state.h" #include "mongo/db/vector_clock.h" -#include "mongo/logv2/log.h" #include "mongo/s/balancer_configuration.h" #include "mongo/s/catalog/type_shard.h" #include "mongo/s/grid.h" @@ -56,30 +54,12 @@ namespace { using ChunkDistributionMap = stdx::unordered_map<ShardId, size_t>; using ZoneShardMap = StringMap<std::vector<ShardId>>; -using boost::intrusive_ptr; - -std::vector<ShardId> getAllNonDrainingShardIdsSorted(OperationContext* opCtx) { - const auto shardsAndOpTime = uassertStatusOKWithContext( - Grid::get(opCtx)->catalogClient()->getAllShards( - opCtx, repl::ReadConcernLevel::kMajorityReadConcern, true /* excludeDraining */), - "Cannot retrieve updated shard list from config server"); - const auto shards = std::move(shardsAndOpTime.value); - const auto lastVisibleOpTime = std::move(shardsAndOpTime.opTime); - - LOGV2_DEBUG(6566600, - 1, - "Successfully retrieved updated shard list from config server", - "nonDrainingShardsNumber"_attr = shards.size(), - "lastVisibleOpTime"_attr = lastVisibleOpTime); - - std::vector<ShardId> shardIds; - std::transform(shards.begin(), - shards.end(), - std::back_inserter(shardIds), - [](const ShardType& shard) { return ShardId(shard.getName()); }); +std::vector<ShardId> getAllShardIdsSorted(OperationContext* opCtx) { + // Many tests assume that chunks will be placed on shards + // according to their IDs in ascending lexical order. + auto shardIds = Grid::get(opCtx)->shardRegistry()->getAllShardIdsNoReload(); std::sort(shardIds.begin(), shardIds.end()); - return shardIds; } @@ -286,8 +266,7 @@ std::unique_ptr<InitialSplitPolicy> InitialSplitPolicy::calculateOptimizationStr const boost::optional<std::vector<BSONObj>>& initialSplitPoints, const std::vector<TagsType>& tags, size_t numShards, - bool collectionIsEmpty, - bool useAutoSplitter) { + bool collectionIsEmpty) { uassert(ErrorCodes::InvalidOptions, str::stream() << "numInitialChunks is only supported when the collection is empty " "and has a hashed field in the shard key pattern", @@ -331,11 +310,7 @@ std::unique_ptr<InitialSplitPolicy> InitialSplitPolicy::calculateOptimizationStr return std::make_unique<SingleChunkOnPrimarySplitPolicy>(); } - if (useAutoSplitter) { - return std::make_unique<AutoSplitInChunksOnPrimaryPolicy>(); - } - - return std::make_unique<SingleChunkOnPrimarySplitPolicy>(); + return std::make_unique<UnoptimizedSplitPolicy>(); } InitialSplitPolicy::ShardCollectionConfig SingleChunkOnPrimarySplitPolicy::createFirstChunks( @@ -359,7 +334,7 @@ InitialSplitPolicy::ShardCollectionConfig SingleChunkOnPrimarySplitPolicy::creat return {std::move(chunks)}; } -InitialSplitPolicy::ShardCollectionConfig AutoSplitInChunksOnPrimaryPolicy::createFirstChunks( +InitialSplitPolicy::ShardCollectionConfig UnoptimizedSplitPolicy::createFirstChunks( OperationContext* opCtx, const ShardKeyPattern& shardKeyPattern, const SplitPolicyParams& params) { @@ -397,7 +372,7 @@ InitialSplitPolicy::ShardCollectionConfig SplitPointsBasedSplitPolicy::createFir const SplitPolicyParams& params) { // On which shards are the generated chunks allowed to be placed. - const auto shardIds = getAllNonDrainingShardIdsSorted(opCtx); + const auto shardIds = getAllShardIdsSorted(opCtx); const auto currentTime = VectorClock::get(opCtx)->getTime(); const auto validAfter = currentTime.clusterTime().asTimestamp(); @@ -428,7 +403,7 @@ InitialSplitPolicy::ShardCollectionConfig AbstractTagsBasedSplitPolicy::createFi const SplitPolicyParams& params) { invariant(!_tags.empty()); - const auto shardIds = getAllNonDrainingShardIdsSorted(opCtx); + const auto shardIds = getAllShardIdsSorted(opCtx); const auto currentTime = VectorClock::get(opCtx)->getTime(); const auto validAfter = currentTime.clusterTime().asTimestamp(); const auto& keyPattern = shardKeyPattern.getKeyPattern(); @@ -675,29 +650,34 @@ std::vector<BSONObj> ReshardingSplitPolicy::createRawPipeline(const ShardKeyPatt std::vector<BSONObj> res; const auto& shardKeyFields = shardKey.getKeyPatternFields(); + + BSONObjBuilder projectValBuilder; BSONObjBuilder sortValBuilder; - using Doc = Document; - using Arr = std::vector<Value>; - using V = Value; - Arr arrayToObjectBuilder; + for (auto&& fieldRef : shardKeyFields) { // If the shard key includes a hashed field and current fieldRef is the hashed field. if (shardKey.isHashedPattern() && fieldRef->dottedField().compare(shardKey.getHashedField().fieldNameStringData()) == 0) { - arrayToObjectBuilder.emplace_back( - Doc{{"k", V{fieldRef->dottedField()}}, - {"v", Doc{{"$toHashedIndexKey", V{"$" + fieldRef->dottedField()}}}}}); + projectValBuilder.append(fieldRef->dottedField(), + BSON("$toHashedIndexKey" + << "$" + fieldRef->dottedField())); } else { - arrayToObjectBuilder.emplace_back(Doc{ - {"k", V{fieldRef->dottedField()}}, - {"v", Doc{{"$ifNull", V{Arr{V{"$" + fieldRef->dottedField()}, V{BSONNULL}}}}}}}); + projectValBuilder.append( + str::stream() << fieldRef->dottedField(), + BSON("$ifNull" << BSON_ARRAY("$" + fieldRef->dottedField() << BSONNULL))); } + sortValBuilder.append(fieldRef->dottedField().toString(), 1); } + + // Do not project _id if it's not part of the shard key. + if (!shardKey.hasId()) { + projectValBuilder.append("_id", 0); + } + res.push_back(BSON("$sample" << BSON("size" << numSplitPoints * samplesPerChunk))); + res.push_back(BSON("$project" << projectValBuilder.obj())); res.push_back(BSON("$sort" << sortValBuilder.obj())); - res.push_back( - Doc{{"$replaceWith", Doc{{"$arrayToObject", Arr{V{arrayToObjectBuilder}}}}}}.toBson()); return res; } @@ -763,12 +743,12 @@ InitialSplitPolicy::ShardCollectionConfig ReshardingSplitPolicy::createFirstChun } { - auto shardIds = getAllNonDrainingShardIdsSorted(opCtx); - for (const auto& shard : shardIds) { + auto allShardIds = getAllShardIdsSorted(opCtx); + for (const auto& shard : allShardIds) { chunkDistribution.emplace(shard, 0); } - zoneToShardMap.emplace("", std::move(shardIds)); + zoneToShardMap.emplace("", std::move(allShardIds)); } std::vector<ChunkType> chunks; @@ -820,34 +800,26 @@ void ReshardingSplitPolicy::_appendSplitPointsFromSample(BSONObjSet* splitPoints while (nextKey && nRemaining > 0) { // if key is hashed, nextKey values are already hashed - auto result = splitPoints->insert(nextKey->getOwned()); + auto result = splitPoints->insert( + dotted_path_support::extractElementsBasedOnTemplate(*nextKey, shardKey.toBSON()) + .getOwned()); + if (result.second) { nRemaining--; } + nextKey = _samples->getNext(); } } std::unique_ptr<ReshardingSplitPolicy::SampleDocumentSource> -ReshardingSplitPolicy::makePipelineDocumentSource_forTest(OperationContext* opCtx, - const NamespaceString& ns, - const ShardKeyPattern& shardKey, - int numInitialChunks, - int samplesPerChunk) { - MakePipelineOptions opts; - opts.attachCursorSource = false; - return _makePipelineDocumentSource( - opCtx, ns, shardKey, numInitialChunks, samplesPerChunk, std::move(opts)); -} - -std::unique_ptr<ReshardingSplitPolicy::SampleDocumentSource> ReshardingSplitPolicy::_makePipelineDocumentSource(OperationContext* opCtx, const NamespaceString& ns, const ShardKeyPattern& shardKey, int numInitialChunks, - int samplesPerChunk, - MakePipelineOptions opts) { + int samplesPerChunk) { auto rawPipeline = createRawPipeline(shardKey, numInitialChunks - 1, samplesPerChunk); + StringMap<ExpressionContext::ResolvedNamespace> resolvedNamespaces; resolvedNamespaces[ns.coll()] = {ns, std::vector<BSONObj>{}}; @@ -861,7 +833,7 @@ ReshardingSplitPolicy::_makePipelineDocumentSource(OperationContext* opCtx, boost::none, /* explain */ false, /* fromMongos */ false, /* needsMerge */ - true, /* allowDiskUse */ + false, /* allowDiskUse */ true, /* bypassDocumentValidation */ false, /* isMapReduceCommand */ ns, @@ -871,10 +843,8 @@ ReshardingSplitPolicy::_makePipelineDocumentSource(OperationContext* opCtx, std::move(resolvedNamespaces), boost::none); /* collUUID */ - expCtx->tempDir = storageGlobalParams.dbpath + "/tmp"; - - return std::make_unique<PipelineDocumentSource>( - Pipeline::makePipeline(rawPipeline, expCtx, opts), samplesPerChunk - 1); + return std::make_unique<PipelineDocumentSource>(Pipeline::makePipeline(rawPipeline, expCtx, {}), + samplesPerChunk - 1); } ReshardingSplitPolicy::PipelineDocumentSource::PipelineDocumentSource( |
