summaryrefslogtreecommitdiff
path: root/src/mongo/db/s/migration_destination_manager.cpp
diff options
context:
space:
mode:
Diffstat (limited to 'src/mongo/db/s/migration_destination_manager.cpp')
-rw-r--r--src/mongo/db/s/migration_destination_manager.cpp162
1 files changed, 44 insertions, 118 deletions
diff --git a/src/mongo/db/s/migration_destination_manager.cpp b/src/mongo/db/s/migration_destination_manager.cpp
index 04399595453..83973741dbe 100644
--- a/src/mongo/db/s/migration_destination_manager.cpp
+++ b/src/mongo/db/s/migration_destination_manager.cpp
@@ -29,6 +29,7 @@
#define MONGO_LOGV2_DEFAULT_COMPONENT ::mongo::logv2::LogComponent::kShardingMigration
+#include "mongo/db/s/migration_batch_fetcher.h"
#include "mongo/platform/basic.h"
#include "mongo/db/s/migration_destination_manager.h"
@@ -39,7 +40,7 @@
#include "mongo/db/auth/authorization_session.h"
#include "mongo/db/cancelable_operation_context.h"
#include "mongo/db/catalog/document_validation.h"
-#include "mongo/db/concurrency/write_conflict_exception.h"
+#include "mongo/db/concurrency/exception_util.h"
#include "mongo/db/db_raii.h"
#include "mongo/db/dbhelpers.h"
#include "mongo/db/index/index_descriptor.h"
@@ -295,6 +296,7 @@ MONGO_FAIL_POINT_DEFINE(migrateThreadHangAtStep3);
MONGO_FAIL_POINT_DEFINE(migrateThreadHangAtStep4);
MONGO_FAIL_POINT_DEFINE(migrateThreadHangAtStep5);
MONGO_FAIL_POINT_DEFINE(migrateThreadHangAtStep6);
+MONGO_FAIL_POINT_DEFINE(migrateThreadHangAfterSteadyTransition);
MONGO_FAIL_POINT_DEFINE(migrateThreadHangAtStep7);
MONGO_FAIL_POINT_DEFINE(failMigrationOnRecipient);
@@ -411,8 +413,8 @@ void MigrationDestinationManager::report(BSONObjBuilder& b,
}
BSONObjBuilder bb(b.subobjStart("counts"));
- bb.append("cloned", _numCloned);
- bb.append("clonedBytes", _clonedBytes);
+ bb.append("cloned", _getNumCloned());
+ bb.append("clonedBytes", _getNumBytesCloned());
bb.append("catchup", _numCatchup);
bb.append("steady", _numSteady);
bb.done();
@@ -446,6 +448,8 @@ Status MigrationDestinationManager::start(OperationContext* opCtx,
_lsid = cloneRequest.getLsid();
_txnNumber = cloneRequest.getTxnNumber();
+ _parallelFetchersSupported = cloneRequest.parallelFetchingSupported();
+
_nss = nss;
_fromShard = cloneRequest.getFromShardId();
_fromShardConnString =
@@ -462,8 +466,8 @@ Status MigrationDestinationManager::start(OperationContext* opCtx,
_chunkMarkedPending = false;
- _numCloned = 0;
- _clonedBytes = 0;
+ _migrationCloningProgress = std::make_shared<MigrationCloningProgressSharedState>();
+
_numCatchup = 0;
_numSteady = 0;
@@ -486,6 +490,9 @@ Status MigrationDestinationManager::start(OperationContext* opCtx,
_sessionMigration = std::make_unique<SessionCatalogMigrationDestination>(
_nss, _fromShard, *_sessionId, _cancellationSource.token());
ShardingStatistics::get(opCtx).countRecipientMoveChunkStarted.addAndFetch(1);
+ if (mongo::feature_flags::gConcurrencyInChunkMigration.isEnabledAndIgnoreFCV())
+ ShardingStatistics::get(opCtx).chunkMigrationConcurrencyCnt.store(
+ chunkMigrationConcurrency.load());
_migrateThreadHandle = stdx::thread([this, cancellationToken = _cancellationSource.token()]() {
_migrateThread(cancellationToken);
@@ -1031,10 +1038,9 @@ void MigrationDestinationManager::cloneCollectionIndexesAndOptions(
<< collectionByUUID->ns());
}
- // We do not have a collection by this name. Create the collection with the donor's
- // options.
+ // We do not have a collection by this name. Create it with the donor's options.
OperationShardingState::ScopedAllowImplicitCollectionCreate_UNSAFE
- unsafeCreateCollection(opCtx);
+ unsafeCreateCollection(opCtx, /* forceCSRAsUnknownAfterCollectionCreation */ true);
WriteUnitOfWork wuow(opCtx);
CollectionOptions collectionOptions = uassertStatusOK(
CollectionOptions::parse(collectionOptionsAndIndexes.options,
@@ -1338,122 +1344,32 @@ void MigrationDestinationManager::_migrateDriver(OperationContext* outerOpCtx,
_sessionMigration->start(opCtx->getServiceContext());
- const BSONObj migrateCloneRequest = createMigrateCloneRequest(_nss, *_sessionId);
-
_chunkMarkedPending = true; // no lock needed, only the migrate thread looks.
- auto assertNotAborted = [&](OperationContext* opCtx) {
- opCtx->checkForInterrupt();
- outerOpCtx->checkForInterrupt();
- uassert(50748, "Migration aborted while copying documents", getState() != kAbort);
- };
-
- auto insertBatchFn = [&](OperationContext* opCtx, BSONObj nextBatch) {
- auto arr = nextBatch["objects"].Obj();
- if (arr.isEmpty()) {
- return false;
- }
- auto it = arr.begin();
- while (it != arr.end()) {
- int batchNumCloned = 0;
- int batchClonedBytes = 0;
- const int batchMaxCloned = migrateCloneInsertionBatchSize.load();
-
- assertNotAborted(opCtx);
-
- write_ops::InsertCommandRequest insertOp(_nss);
- insertOp.getWriteCommandRequestBase().setOrdered(true);
- insertOp.setDocuments([&] {
- std::vector<BSONObj> toInsert;
- while (it != arr.end() &&
- (batchMaxCloned <= 0 || batchNumCloned < batchMaxCloned)) {
- const auto& doc = *it;
- BSONObj docToClone = doc.Obj();
- toInsert.push_back(docToClone);
- batchNumCloned++;
- batchClonedBytes += docToClone.objsize();
- ++it;
- }
- return toInsert;
- }());
-
- {
- // Disable the schema validation (during document inserts and updates)
- // and any internal validation for opCtx for performInserts()
- DisableDocumentValidation documentValidationDisabler(
- opCtx,
- DocumentValidationSettings::kDisableSchemaValidation |
- DocumentValidationSettings::kDisableInternalValidation);
- const auto reply = write_ops_exec::performInserts(
- opCtx, insertOp, OperationSource::kFromMigrate);
- for (unsigned long i = 0; i < reply.results.size(); ++i) {
- uassertStatusOKWithContext(reply.results[i],
- str::stream() << "Insert of "
- << insertOp.getDocuments()[i]
- << " failed.");
- }
- // Revert to the original DocumentValidationSettings for opCtx
- }
-
- migrationutil::persistUpdatedNumOrphans(
- opCtx, _migrationId.get(), *_collectionUuid, batchNumCloned);
-
- {
- stdx::lock_guard<Latch> statsLock(_mutex);
- _numCloned += batchNumCloned;
- ShardingStatistics::get(opCtx).countDocsClonedOnRecipient.addAndFetch(
- batchNumCloned);
- _clonedBytes += batchClonedBytes;
- }
- if (_writeConcern.needToWaitForOtherNodes()) {
- runWithoutSession(outerOpCtx, [&] {
- repl::ReplicationCoordinator::StatusAndDuration replStatus =
- repl::ReplicationCoordinator::get(opCtx)->awaitReplication(
- opCtx,
- repl::ReplClientInfo::forClient(opCtx->getClient()).getLastOp(),
- _writeConcern);
- if (replStatus.status.code() == ErrorCodes::WriteConcernFailed) {
- LOGV2_WARNING(
- 22011,
- "secondaryThrottle on, but doc insert timed out; continuing",
- "migrationId"_attr = _migrationId->toBSON());
- } else {
- uassertStatusOK(replStatus.status);
- }
- });
- }
-
- sleepmillis(migrateCloneInsertionBatchDelayMS.load());
- }
- return true;
- };
-
- auto fetchBatchFn = [&](OperationContext* opCtx, BSONObj* nextBatch) {
- auto commandResponse = uassertStatusOKWithContext(
- fromShard->runCommand(opCtx,
- ReadPreferenceSetting(ReadPreference::PrimaryOnly),
- "admin",
- migrateCloneRequest,
- Shard::RetryPolicy::kNoRetry),
- "_migrateClone failed: ");
-
- uassertStatusOKWithContext(
- Shard::CommandResponse::getEffectiveStatus(commandResponse),
- "_migrateClone failed: ");
-
- *nextBatch = commandResponse.response;
- return nextBatch->getField("objects").Obj().isEmpty();
- };
-
- // If running on a replicated system, we'll need to flush the docs we cloned to the
- // secondaries
- lastOpApplied = fetchAndApplyBatch(opCtx, insertBatchFn, fetchBatchFn);
+ {
+ // Destructor of MigrationBatchFetcher is non-trivial. Therefore,
+ // this scope has semantic significance.
+ MigrationBatchFetcher<MigrationBatchInserter> fetcher{outerOpCtx,
+ opCtx,
+ _nss,
+ *_sessionId,
+ _writeConcern,
+ _fromShard,
+ range,
+ *_migrationId,
+ *_collectionUuid,
+ _migrationCloningProgress,
+ _parallelFetchersSupported};
+ fetcher.fetchAndScheduleInsertion();
+ }
+ opCtx->checkForInterrupt();
+ lastOpApplied = _migrationCloningProgress->getMaxOptime();
timing->done(4);
migrateThreadHangAtStep4.pauseWhileSet();
if (MONGO_unlikely(failMigrationOnRecipient.shouldFail())) {
- _setStateFail(str::stream() << "failing migration after cloning " << _numCloned
+ _setStateFail(str::stream() << "failing migration after cloning " << _getNumCloned()
<< " docs due to failMigrationOnRecipient failpoint");
return;
}
@@ -1491,6 +1407,8 @@ void MigrationDestinationManager::_migrateDriver(OperationContext* outerOpCtx,
if (!_applyMigrateOp(opCtx, nextBatch)) {
return true;
}
+ ShardingStatistics::get(opCtx).countBytesClonedOnCatchUpOnRecipient.addAndFetch(
+ nextBatch["size"].number());
const int maxIterations = 3600 * 50;
@@ -1569,6 +1487,7 @@ void MigrationDestinationManager::_migrateDriver(OperationContext* outerOpCtx,
{
// 6. Wait for commit
_setState(kSteady);
+ migrateThreadHangAfterSteadyTransition.pauseWhileSet();
bool transferAfterCommit = false;
while (getState() == kSteady || getState() == kCommitStart) {
@@ -1596,7 +1515,8 @@ void MigrationDestinationManager::_migrateDriver(OperationContext* outerOpCtx,
auto mods = res.response;
- if (mods["size"].number() > 0 && _applyMigrateOp(opCtx, mods)) {
+ if (mods["size"].number() > 0) {
+ (void)_applyMigrateOp(opCtx, mods);
lastOpApplied = repl::ReplClientInfo::forClient(opCtx->getClient()).getLastOp();
continue;
}
@@ -1755,6 +1675,7 @@ void MigrationDestinationManager::_migrateDriver(OperationContext* outerOpCtx,
bool MigrationDestinationManager::_applyMigrateOp(OperationContext* opCtx, const BSONObj& xfer) {
bool didAnything = false;
long long changeInOrphans = 0;
+ long long totalDocs = 0;
// Deleted documents
if (xfer["deleted"].isABSONObj()) {
@@ -1765,6 +1686,7 @@ bool MigrationDestinationManager::_applyMigrateOp(OperationContext* opCtx, const
BSONObjIterator i(xfer["deleted"].Obj());
while (i.more()) {
+ totalDocs++;
AutoGetCollection autoColl(opCtx, _nss, MODE_IX);
uassert(ErrorCodes::ConflictingOperationInProgress,
str::stream() << "Collection " << _nss.ns()
@@ -1807,6 +1729,7 @@ bool MigrationDestinationManager::_applyMigrateOp(OperationContext* opCtx, const
if (xfer["reload"].isABSONObj()) {
BSONObjIterator i(xfer["reload"].Obj());
while (i.more()) {
+ totalDocs++;
AutoGetCollection autoColl(opCtx, _nss, MODE_IX);
uassert(ErrorCodes::ConflictingOperationInProgress,
str::stream() << "Collection " << _nss.ns()
@@ -1861,6 +1784,9 @@ bool MigrationDestinationManager::_applyMigrateOp(OperationContext* opCtx, const
migrationutil::persistUpdatedNumOrphans(
opCtx, _migrationId.get(), *_collectionUuid, changeInOrphans);
}
+
+ ShardingStatistics::get(opCtx).countDocsClonedOnCatchUpOnRecipient.addAndFetch(totalDocs);
+
return didAnything;
}