summaryrefslogtreecommitdiff
path: root/src/mongo/db/s/config/initial_split_policy.cpp
diff options
context:
space:
mode:
Diffstat (limited to 'src/mongo/db/s/config/initial_split_policy.cpp')
-rw-r--r--src/mongo/db/s/config/initial_split_policy.cpp108
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(