diff options
5 files changed, 145 insertions, 3 deletions
diff --git a/jstests/sharding/block_chunk_migrations_without_hashed_shard_key_index.js b/jstests/sharding/block_chunk_migrations_without_hashed_shard_key_index.js new file mode 100644 index 00000000000..e03cb160c89 --- /dev/null +++ b/jstests/sharding/block_chunk_migrations_without_hashed_shard_key_index.js @@ -0,0 +1,75 @@ +/** + * Tests that chunk migrations are blocked when there is no index on a hashed shard key. + * + * @tags: [ + * requires_fcv_70, + * featureFlagShardKeyIndexOptionalHashedSharding + * ] + */ + +(function() { +"use strict"; + +load("jstests/sharding/libs/find_chunks_util.js"); + +const st = new ShardingTest({}); + +const dbName = "testDb"; +const collName = "testColl"; +const nss = dbName + "." + collName; +const configDB = st.s.getDB('config'); + +const coll = st.getDB(dbName).getCollection(collName); +const kDataString = "a".repeat(1024 * 1024); +let docs = Array.from({length: 1000}, (x, i) => ({_id: i, field: kDataString})); + +assert.commandWorked( + st.s.adminCommand({enablesharding: dbName, primaryShard: st.shard0.shardName})); +assert.commandWorked(coll.createIndex({"_id": "hashed"})); +assert.commandWorked(st.s.adminCommand({shardCollection: nss, key: {_id: "hashed"}})); + +// Move all chunks to a single shard so the balancer is triggered due to data imbalance. +let chunks = findChunksUtil.findChunksByNs(configDB, nss).toArray(); +chunks.forEach(chunk => { + st.s.adminCommand({moveChunk: nss, bounds: [chunk.min, chunk.max], to: st.shard0.shardName}); +}); + +assert.eq(0, findChunksUtil.findChunksByNs(configDB, nss, {shard: st.shard1.shardName}).itcount()); + +assert.commandWorked(coll.insert(docs)); +assert.commandWorked(coll.dropIndex({"_id": "hashed"})); + +st.startBalancer(); +st.awaitBalancerRound(); + +// During balancing, the balancer should catch the IndexNotFound and turn off the balancer for the +// collection by setting {noBalance : true}. +assert.soon(() => { + return configDB.getCollection('collections').findOne({_id: nss}).noBalance === true; +}); + +// Confirm all chunks remain on shard0. +assert.eq(0, findChunksUtil.findChunksByNs(configDB, nss, {shard: st.shard1.shardName}).itcount()); + +// Commands that trigger chunk migrations should fail with IndexNotFound. +assert.commandFailedWithCode( + st.s.adminCommand( + {moveChunk: nss, bounds: [chunks[0].min, chunks[0].max], to: st.shard1.shardName}), + ErrorCodes.IndexNotFound); +assert.commandFailedWithCode( + st.s.adminCommand( + {moveRange: nss, toShard: st.shard1.shardName, min: chunks[0].min, max: chunks[0].max}), + ErrorCodes.IndexNotFound); + +// Recreate the index and verify that we can re-enable balancing. +assert.commandWorked(coll.createIndex({"_id": "hashed"})); +st.enableBalancing(nss); + +assert.eq(false, configDB.getCollection('collections').findOne({_id: nss}).noBalance); +st.awaitBalancerRound(); +assert.soon(() => { + return findChunksUtil.findChunksByNs(configDB, nss, {shard: st.shard1.shardName}).itcount() > 0; +}); + +st.stop(); +})(); diff --git a/src/mongo/db/s/balancer/balancer.cpp b/src/mongo/db/s/balancer/balancer.cpp index d92bb2a4b2a..77d8af69baa 100644 --- a/src/mongo/db/s/balancer/balancer.cpp +++ b/src/mongo/db/s/balancer/balancer.cpp @@ -50,6 +50,7 @@ #include "mongo/db/s/config/sharding_catalog_manager.h" #include "mongo/db/s/sharding_config_server_parameters_gen.h" #include "mongo/db/s/sharding_logging.h" +#include "mongo/db/server_feature_flags_gen.h" #include "mongo/executor/scoped_task_executor.h" #include "mongo/logv2/log.h" #include "mongo/s/balancer_configuration.h" @@ -1138,6 +1139,28 @@ int Balancer::_moveChunks(OperationContext* opCtx, continue; } + if (status == ErrorCodes::IndexNotFound && + gFeatureFlagShardKeyIndexOptionalHashedSharding.isEnabled( + serverGlobalParams.featureCompatibility)) { + + const auto [cm, _] = uassertStatusOK( + Grid::get(opCtx)->catalogCache()->getCollectionRoutingInfoWithRefresh( + opCtx, migrateInfo.nss)); + + if (cm.getShardKeyPattern().isHashedPattern()) { + LOGV2(78252, + "Turning off balancing for hashed collection because migration failed due to " + "missing shardkey index", + "migrateInfo"_attr = redact(migrateInfo.toString()), + "error"_attr = redact(status), + "collection"_attr = migrateInfo.nss); + + // Schedule writing to config.collections to turn off the balancer. + _commandScheduler->disableBalancerForCollection(opCtx, migrateInfo.nss); + continue; + } + } + LOGV2(21872, "Migration {migrateInfo} failed with {error}", "Migration failed", diff --git a/src/mongo/db/s/balancer/balancer_commands_scheduler.h b/src/mongo/db/s/balancer/balancer_commands_scheduler.h index 2c0e69000c0..d308758c222 100644 --- a/src/mongo/db/s/balancer/balancer_commands_scheduler.h +++ b/src/mongo/db/s/balancer/balancer_commands_scheduler.h @@ -62,6 +62,9 @@ public: */ virtual void stop() = 0; + virtual void disableBalancerForCollection(OperationContext* opCtx, + const NamespaceString& nss) = 0; + virtual SemiFuture<void> requestMergeChunks(OperationContext* opCtx, const NamespaceString& nss, const ShardId& shardId, diff --git a/src/mongo/db/s/balancer/balancer_commands_scheduler_impl.cpp b/src/mongo/db/s/balancer/balancer_commands_scheduler_impl.cpp index 82db31f0437..2130f88d2fb 100644 --- a/src/mongo/db/s/balancer/balancer_commands_scheduler_impl.cpp +++ b/src/mongo/db/s/balancer/balancer_commands_scheduler_impl.cpp @@ -148,6 +148,17 @@ void BalancerCommandsSchedulerImpl::stop() { _workerThreadHandle.join(); } +void BalancerCommandsSchedulerImpl::disableBalancerForCollection(OperationContext* opCtx, + const NamespaceString& nss) { + auto commandInfo = std::make_shared<DisableBalancerCommandInfo>(nss, ShardId::kConfigServerId); + + _buildAndEnqueueNewRequest(opCtx, std::move(commandInfo)) + .then([](const executor::RemoteCommandResponse& remoteResponse) { + return processRemoteResponse(remoteResponse); + }) + .getAsync([](auto) {}); +} + SemiFuture<void> BalancerCommandsSchedulerImpl::requestMoveRange( OperationContext* opCtx, const ShardsvrMoveRange& request, diff --git a/src/mongo/db/s/balancer/balancer_commands_scheduler_impl.h b/src/mongo/db/s/balancer/balancer_commands_scheduler_impl.h index 627d3a08348..8949bde504e 100644 --- a/src/mongo/db/s/balancer/balancer_commands_scheduler_impl.h +++ b/src/mongo/db/s/balancer/balancer_commands_scheduler_impl.h @@ -36,6 +36,7 @@ #include "mongo/executor/scoped_task_executor.h" #include "mongo/platform/mutex.h" #include "mongo/rpc/get_status_from_command_result.h" +#include "mongo/s/catalog/type_collection.h" #include "mongo/s/client/shard.h" #include "mongo/s/request_types/merge_chunk_request_gen.h" #include "mongo/s/request_types/migration_secondary_throttle_options.h" @@ -111,6 +112,10 @@ private: boost::optional<ExternalClientInfo> _clientInfo; }; +/** + * Set of command-specific subclasses of CommandInfo. + */ + class MoveRangeCommandInfo : public CommandInfo { public: MoveRangeCommandInfo(const ShardsvrMoveRange& request, @@ -137,9 +142,32 @@ private: const WriteConcernOptions _wc; }; -/** - * Set of command-specific subclasses of CommandInfo. - */ +class DisableBalancerCommandInfo : public CommandInfo { +public: + DisableBalancerCommandInfo(const NamespaceString& nss, const ShardId& shardId) + : CommandInfo(shardId, nss, boost::none) {} + + BSONObj serialise() const override { + BSONObjBuilder updateCmd; + updateCmd.append("$set", BSON("noBalance" << true)); + + const auto updateOp = BatchedCommandRequest::buildUpdateOp( + CollectionType::ConfigNS, + BSON(CollectionType::kNssFieldName + << NamespaceStringUtil::serialize(getNameSpace())) /* query */, + updateCmd.obj() /* update */, + false /* upsert */, + false /* multi */); + BSONObjBuilder cmdObj(updateOp.toBSON()); + cmdObj.append(WriteConcernOptions::kWriteConcernField, + WriteConcernOptions::kInternalWriteDefault); + return cmdObj.obj(); + } + + std::string getTargetDb() const override { + return DatabaseName::kConfig.toString(); + } +}; class MergeChunksCommandInfo : public CommandInfo { public: @@ -361,6 +389,8 @@ public: void stop() override; + void disableBalancerForCollection(OperationContext* opCtx, const NamespaceString& nss) override; + SemiFuture<void> requestMoveRange(OperationContext* opCtx, const ShardsvrMoveRange& request, const WriteConcernOptions& secondaryThrottleWC, |
