summaryrefslogtreecommitdiff
path: root/src/mongo/db/commands/write_commands.cpp
diff options
context:
space:
mode:
Diffstat (limited to 'src/mongo/db/commands/write_commands.cpp')
-rw-r--r--src/mongo/db/commands/write_commands.cpp207
1 files changed, 95 insertions, 112 deletions
diff --git a/src/mongo/db/commands/write_commands.cpp b/src/mongo/db/commands/write_commands.cpp
index 0254baca47d..9186ff13cbe 100644
--- a/src/mongo/db/commands/write_commands.cpp
+++ b/src/mongo/db/commands/write_commands.cpp
@@ -30,6 +30,7 @@
#define MONGO_LOGV2_DEFAULT_COMPONENT ::mongo::logv2::LogComponent::kDefault
#include "mongo/base/checked_cast.h"
+#include "mongo/base/error_codes.h"
#include "mongo/bson/bsonobjbuilder.h"
#include "mongo/bson/mutable/document.h"
#include "mongo/bson/mutable/element.h"
@@ -73,6 +74,7 @@
#include "mongo/db/timeseries/bucket_catalog.h"
#include "mongo/db/timeseries/bucket_compression.h"
#include "mongo/db/timeseries/timeseries_constants.h"
+#include "mongo/db/timeseries/timeseries_extended_range.h"
#include "mongo/db/timeseries/timeseries_options.h"
#include "mongo/db/timeseries/timeseries_stats.h"
#include "mongo/db/transaction_participant.h"
@@ -527,6 +529,11 @@ public:
}
write_ops::InsertCommandReply typedRun(OperationContext* opCtx) final try {
+ // On debug builds, verify that the estimated size of the insert command is at least as
+ // large as the size of the actual, serialized insert command. This ensures that the
+ // logic which estimates the size of insert commands is correct.
+ dassert(write_ops::verifySizeEstimate(request(), &unparsedRequest()));
+
transactionChecks(opCtx, ns());
if (request().getEncryptionInformation().has_value() &&
@@ -698,16 +705,16 @@ public:
OperationSource::kTimeseriesInsert));
}
- TimeseriesSingleWriteResult _performTimeseriesBucketCompression(
+ void _performTimeseriesBucketCompression(
OperationContext* opCtx, const BucketCatalog::ClosedBucket& closedBucket) const {
if (!feature_flags::gTimeseriesBucketCompression.isEnabled(
serverGlobalParams.featureCompatibility)) {
- return {SingleWriteResult(), true};
+ return;
}
// Buckets with just a single measurement is not worth compressing.
if (closedBucket.numMeasurements <= 1) {
- return {SingleWriteResult(), true};
+ return;
}
bool validateCompression = gValidateTimeseriesCompression.load();
@@ -745,8 +752,8 @@ public:
auto compressionOp =
_makeTimeseriesCompressionOp(opCtx, closedBucket.bucketId, bucketCompressionFunc);
- auto result = _getTimeseriesSingleWriteResult(
- write_ops_exec::performUpdates(opCtx, compressionOp, OperationSource::kStandard));
+ auto result = _getTimeseriesSingleWriteResult(write_ops_exec::performUpdates(
+ opCtx, compressionOp, OperationSource::kTimeseriesBucketCompression));
// Report stats, if we fail before running the transform function then just skip
// reporting.
@@ -761,8 +768,6 @@ public:
stats.onBucketClosed(*beforeSize, compressionStats);
}
}
-
- return result;
}
/**
@@ -776,7 +781,8 @@ public:
std::vector<write_ops::WriteError>* errors,
boost::optional<repl::OpTime>* opTime,
boost::optional<OID>* electionId,
- std::vector<size_t>* docsToRetry) const try {
+ std::vector<size_t>* docsToRetry,
+ absl::flat_hash_map<int, int>& retryAttemptsForDup) const try {
auto& bucketCatalog = BucketCatalog::get(opCtx);
auto metadata = bucketCatalog.getMetadata(batch->bucket());
@@ -796,9 +802,18 @@ public:
_performTimeseriesInsert(opCtx, batch, metadata, std::move(stmtIds));
if (auto error =
generateError(opCtx, output.result, start + index, errors->size())) {
- errors->emplace_back(std::move(*error));
- bucketCatalog.abort(batch, output.result.getStatus());
- return output.canContinue;
+ bool canContinue = output.canContinue;
+ // Automatically attempts to retry on DuplicateKey error.
+ if (error->getStatus().code() == ErrorCodes::DuplicateKey &&
+ retryAttemptsForDup[index]++ <
+ gTimeseriesInsertMaxRetriesOnDuplicates.load()) {
+ docsToRetry->push_back(index);
+ canContinue = true;
+ } else {
+ errors->emplace_back(std::move(*error));
+ }
+ BucketCatalog::get(opCtx).abort(batch, output.result.getStatus());
+ return canContinue;
}
invariant(output.result.getValue().getN() == 1,
@@ -828,12 +843,7 @@ public:
if (closedBucket) {
// If this write closed a bucket, compress the bucket
- auto output = _performTimeseriesBucketCompression(opCtx, *closedBucket);
- if (auto error =
- generateError(opCtx, output.result, start + index, errors->size())) {
- errors->emplace_back(std::move(*error));
- return output.canContinue;
- }
+ _performTimeseriesBucketCompression(opCtx, *closedBucket);
}
return true;
} catch (const DBException& ex) {
@@ -841,19 +851,12 @@ public:
throw;
}
- enum struct TimeseriesAtomicWriteResult {
- kSuccess,
- kContinuableError,
- kNonContinuableError,
- };
-
- TimeseriesAtomicWriteResult _commitTimeseriesBucketsAtomically(
- OperationContext* opCtx,
- TimeseriesBatches* batches,
- TimeseriesStmtIds&& stmtIds,
- std::vector<write_ops::WriteError>* errors,
- boost::optional<repl::OpTime>* opTime,
- boost::optional<OID>* electionId) const {
+ bool _commitTimeseriesBucketsAtomically(OperationContext* opCtx,
+ TimeseriesBatches* batches,
+ TimeseriesStmtIds&& stmtIds,
+ std::vector<write_ops::WriteError>* errors,
+ boost::optional<repl::OpTime>* opTime,
+ boost::optional<OID>* electionId) const {
auto& bucketCatalog = BucketCatalog::get(opCtx);
std::vector<std::reference_wrapper<std::shared_ptr<BucketCatalog::WriteBatch>>>
@@ -866,7 +869,7 @@ public:
}
if (batchesToCommit.empty()) {
- return TimeseriesAtomicWriteResult::kSuccess;
+ return true;
}
// Sort by bucket so that preparing the commit for each batch cannot deadlock.
@@ -892,7 +895,7 @@ public:
auto prepareCommitStatus = bucketCatalog.prepareCommit(batch);
if (!prepareCommitStatus.isOK()) {
abortStatus = prepareCommitStatus;
- return TimeseriesAtomicWriteResult::kContinuableError;
+ return false;
}
if (batch.get()->numPreviouslyCommittedMeasurements() == 0) {
@@ -909,33 +912,26 @@ public:
auto result =
write_ops_exec::performAtomicTimeseriesWrites(opCtx, insertOps, updateOps);
if (!result.isOK()) {
+ if (result.code() == ErrorCodes::DuplicateKey) {
+ BucketCatalog::get(opCtx).resetBucketOIDCounter();
+ }
abortStatus = result;
- return TimeseriesAtomicWriteResult::kContinuableError;
+ return false;
}
getOpTimeAndElectionId(opCtx, opTime, electionId);
- bool compressClosedBuckets = true;
for (auto batch : batchesToCommit) {
auto closedBucket = bucketCatalog.finish(
batch, BucketCatalog::CommitInfo{*opTime, *electionId});
batch.get().reset();
- if (!closedBucket || !compressClosedBuckets) {
+ if (!closedBucket) {
continue;
}
// If this write closed a bucket, compress the bucket
- auto ret = _performTimeseriesBucketCompression(opCtx, *closedBucket);
- if (!ret.result.isOK()) {
- // Don't try to compress any other buckets if we fail. We're not allowed to
- // do more write operations.
- compressClosedBuckets = false;
- }
- if (!ret.canContinue) {
- abortStatus = ret.result.getStatus();
- return TimeseriesAtomicWriteResult::kNonContinuableError;
- }
+ _performTimeseriesBucketCompression(opCtx, *closedBucket);
}
} catch (const DBException& ex) {
abortStatus = ex.toStatus();
@@ -943,7 +939,7 @@ public:
}
batchGuard.dismiss();
- return TimeseriesAtomicWriteResult::kSuccess;
+ return true;
}
// For sharded time-series collections, we need to use the granularity from the config
@@ -968,10 +964,7 @@ public:
}
}
- std::tuple<TimeseriesBatches,
- TimeseriesStmtIds,
- size_t /* numInserted */,
- bool /* canContinue */>
+ std::tuple<TimeseriesBatches, TimeseriesStmtIds, size_t /* numInserted */>
_insertIntoBucketCatalog(OperationContext* opCtx,
size_t start,
size_t numDocs,
@@ -1011,7 +1004,6 @@ public:
TimeseriesBatches batches;
TimeseriesStmtIds stmtIds;
- bool canContinue = true;
auto insert = [&](size_t index) {
invariant(start + index < request().getDocuments().size());
@@ -1057,22 +1049,8 @@ public:
// If this insert closed buckets, rewrite to be a compressed column. If we cannot
// perform write operations at this point the bucket will be left uncompressed.
for (const auto& closedBucket : result.getValue().closedBuckets) {
- if (!canContinue) {
- break;
- }
-
// If this write closed a bucket, compress the bucket
- auto ret = _performTimeseriesBucketCompression(opCtx, closedBucket);
- if (auto error =
- generateError(opCtx, ret.result, start + index, errors->size())) {
- // Bucket compression only fail when we may not try to perform any other
- // write operation. When handleError() inside write_ops_exec.cpp return
- // false.
- errors->emplace_back(std::move(*error));
- canContinue = false;
- return false;
- }
- canContinue = ret.canContinue;
+ _performTimeseriesBucketCompression(opCtx, closedBucket);
}
return true;
@@ -1083,15 +1061,12 @@ public:
} else {
for (size_t i = 0; i < numDocs; i++) {
if (!insert(i) && request().getOrdered()) {
- return {std::move(batches), std::move(stmtIds), i, canContinue};
+ return {std::move(batches), std::move(stmtIds), i};
}
}
}
- return {std::move(batches),
- std::move(stmtIds),
- request().getDocuments().size(),
- canContinue};
+ return {std::move(batches), std::move(stmtIds), request().getDocuments().size()};
}
void _getTimeseriesBatchResults(OperationContext* opCtx,
@@ -1153,30 +1128,25 @@ public:
}
}
- TimeseriesAtomicWriteResult _performOrderedTimeseriesWritesAtomically(
- OperationContext* opCtx,
- std::vector<write_ops::WriteError>* errors,
- boost::optional<repl::OpTime>* opTime,
- boost::optional<OID>* electionId,
- bool* containsRetry) const {
- auto [batches, stmtIds, numInserted, canContinue] = _insertIntoBucketCatalog(
+ bool _performOrderedTimeseriesWritesAtomically(OperationContext* opCtx,
+ std::vector<write_ops::WriteError>* errors,
+ boost::optional<repl::OpTime>* opTime,
+ boost::optional<OID>* electionId,
+ bool* containsRetry) const {
+ auto [batches, stmtIds, numInserted] = _insertIntoBucketCatalog(
opCtx, 0, request().getDocuments().size(), {}, errors, containsRetry);
- if (!canContinue) {
- return TimeseriesAtomicWriteResult::kNonContinuableError;
- }
hangTimeseriesInsertBeforeCommit.pauseWhileSet();
- auto result = _commitTimeseriesBucketsAtomically(
- opCtx, &batches, std::move(stmtIds), errors, opTime, electionId);
- if (result != TimeseriesAtomicWriteResult::kSuccess) {
- return result;
+ if (!_commitTimeseriesBucketsAtomically(
+ opCtx, &batches, std::move(stmtIds), errors, opTime, electionId)) {
+ return false;
}
_getTimeseriesBatchResults(
opCtx, batches, 0, batches.size(), true, errors, opTime, electionId);
- return TimeseriesAtomicWriteResult::kSuccess;
+ return true;
}
/**
@@ -1187,19 +1157,9 @@ public:
boost::optional<repl::OpTime>* opTime,
boost::optional<OID>* electionId,
bool* containsRetry) const {
- auto result = _performOrderedTimeseriesWritesAtomically(
- opCtx, errors, opTime, electionId, containsRetry);
- switch (result) {
- case TimeseriesAtomicWriteResult::kSuccess:
- return request().getDocuments().size();
- case TimeseriesAtomicWriteResult::kNonContinuableError:
- // If we can't continue, we know that 0 were inserted since this function should
- // guarantee that the inserts are atomic.
- return 0;
- case TimeseriesAtomicWriteResult::kContinuableError:
- break;
- default:
- MONGO_UNREACHABLE;
+ if (_performOrderedTimeseriesWritesAtomically(
+ opCtx, errors, opTime, electionId, containsRetry)) {
+ return request().getDocuments().size();
}
for (size_t i = 0; i < request().getDocuments().size(); ++i) {
@@ -1218,6 +1178,8 @@ public:
* which were attempted in an update operation, but found no bucket to update. These indices
* can be passed as the 'indices' parameter in a subsequent call to this function, in order
* to to be retried.
+ * In rare cases due to collision from OID generation, we will also retry inserting those
+ * bucket * documents for a limited number of times.
*/
std::vector<size_t> _performUnorderedTimeseriesWrites(
OperationContext* opCtx,
@@ -1227,17 +1189,16 @@ public:
std::vector<write_ops::WriteError>* errors,
boost::optional<repl::OpTime>* opTime,
boost::optional<OID>* electionId,
- bool* containsRetry) const {
- auto [batches, bucketStmtIds, _, canContinue] =
+ bool* containsRetry,
+ absl::flat_hash_map<int, int>& retryAttemptsForDup) const {
+ auto [batches, bucketStmtIds, _] =
_insertIntoBucketCatalog(opCtx, start, numDocs, indices, errors, containsRetry);
hangTimeseriesInsertBeforeCommit.pauseWhileSet();
std::vector<size_t> docsToRetry;
- if (!canContinue) {
- return docsToRetry;
- }
+ bool canContinue = true;
size_t itr = 0;
for (; itr < batches.size(); ++itr) {
@@ -1246,7 +1207,6 @@ public:
auto stmtIds = isTimeseriesWriteRetryable(opCtx)
? std::move(bucketStmtIds[batch->bucket().id])
: std::vector<StmtId>{};
-
canContinue = _commitTimeseriesBucket(opCtx,
batch,
start,
@@ -1255,7 +1215,8 @@ public:
errors,
opTime,
electionId,
- &docsToRetry);
+ &docsToRetry,
+ retryAttemptsForDup);
batch.reset();
if (!canContinue) {
break;
@@ -1280,9 +1241,20 @@ public:
boost::optional<OID>* electionId,
bool* containsRetry) const {
std::vector<size_t> docsToRetry;
+ absl::flat_hash_map<int, int> retryAttemptsForDup;
do {
- docsToRetry = _performUnorderedTimeseriesWrites(
- opCtx, start, numDocs, docsToRetry, errors, opTime, electionId, containsRetry);
+ docsToRetry = _performUnorderedTimeseriesWrites(opCtx,
+ start,
+ numDocs,
+ docsToRetry,
+ errors,
+ opTime,
+ electionId,
+ containsRetry,
+ retryAttemptsForDup);
+ if (!retryAttemptsForDup.empty()) {
+ BucketCatalog::get(opCtx).resetBucketOIDCounter();
+ }
} while (!docsToRetry.empty());
}
@@ -1441,19 +1413,25 @@ public:
invariant(!_commandObj.isEmpty());
- if (const auto& shardVersion = _commandObj.getField("shardVersion");
- !shardVersion.eoo()) {
- bob->append(shardVersion);
- }
bob->append("find", _commandObj["update"].String());
extractQueryDetails(_updateOpObj, bob);
bob->append("batchSize", 1);
bob->append("singleBatch", true);
+
+ if (const auto& shardVersion = _commandObj.getField("shardVersion");
+ !shardVersion.eoo()) {
+ bob->append(shardVersion);
+ }
}
write_ops::UpdateCommandReply typedRun(OperationContext* opCtx) final try {
+ // On debug builds, verify that the estimated size of the update command is at least as
+ // large as the size of the actual, serialized update command. This ensures that the
+ // logic which estimates the size of update commands is correct.
+ dassert(write_ops::verifySizeEstimate(request(), &unparsedRequest()));
transactionChecks(opCtx, ns());
+
write_ops::UpdateCommandReply updateReply;
OperationSource source = OperationSource::kStandard;
@@ -1641,6 +1619,11 @@ public:
}
write_ops::DeleteCommandReply typedRun(OperationContext* opCtx) final try {
+ // On debug builds, verify that the estimated size of the deletes are at least as large
+ // as the actual, serialized size. This ensures that the logic that estimates the size
+ // of deletes for batch writes is correct.
+ dassert(write_ops::verifySizeEstimate(request(), &unparsedRequest()));
+
transactionChecks(opCtx, ns());
write_ops::DeleteCommandReply deleteReply;
OperationSource source = OperationSource::kStandard;