diff options
| author | Lucas de Castro Borges <lucas@gnuabordo.com.br> | 2025-02-14 14:26:38 -0300 |
|---|---|---|
| committer | Lucas de Castro Borges <lucas@gnuabordo.com.br> | 2025-02-14 14:26:38 -0300 |
| commit | 294bc6ecabf14c09c9bc8644704921dcf97cb44e (patch) | |
| tree | 279b1e0bab53901a1647ac63c1c724f0f789a663 /src/mongo/db/pipeline/document_source_writer.h | |
| parent | 70be7c27a251621187a1de533462ae2bb1e3bd39 (diff) | |
| parent | 1e917fd798aa25b7066d4b414b51184f13d5a092 (diff) | |
Update upstream source from tag 'upstream/6.0.10'debian/6.0.10-1
Update to upstream version '6.0.10'
with Debian dir 2d176fa254eee97b139f712fec5709641335a8c3
Diffstat (limited to 'src/mongo/db/pipeline/document_source_writer.h')
| -rw-r--r-- | src/mongo/db/pipeline/document_source_writer.h | 93 |
1 files changed, 80 insertions, 13 deletions
diff --git a/src/mongo/db/pipeline/document_source_writer.h b/src/mongo/db/pipeline/document_source_writer.h index a94b4efca6b..25fd08aac5c 100644 --- a/src/mongo/db/pipeline/document_source_writer.h +++ b/src/mongo/db/pipeline/document_source_writer.h @@ -36,8 +36,11 @@ #include "mongo/db/db_raii.h" #include "mongo/db/operation_context.h" #include "mongo/db/pipeline/document_source.h" +#include "mongo/db/query/query_knobs_gen.h" #include "mongo/db/read_concern.h" #include "mongo/db/storage/recovery_unit.h" +#include "mongo/rpc/metadata/impersonated_user_metadata.h" +#include "mongo/s/write_ops/batched_command_request.h" namespace mongo { using namespace fmt::literals; @@ -83,11 +86,14 @@ public: /** * This is a base abstract class for all stages performing a write operation into an output * collection. The writes are organized in batches in which elements are objects of the templated - * type 'B'. A subclass must override two methods to be able to write into the output collection: + * type 'B'. A subclass must override the following methods to be able to write into the output + * collection: * - * 1. 'makeBatchObject()' - to create an object of type 'B' from the given 'Document', which is, + * - 'makeBatchObject()' - creates an object of type 'B' from the given 'Document', which is, * essentially, a result of the input source's 'getNext()' . - * 2. 'spill()' - to write the batch into the output collection. + * - 'spill()' - writes the batch into the output collection. + * - 'initializeBatchedWriteRequest()' - initializes the request object for writing a batch to + * the output collection. * * Two other virtual methods exist which a subclass may override: 'initialize()' and 'finalize()', * which are called before the first element is read from the input source, and after the last one @@ -99,12 +105,26 @@ public: using BatchObject = B; using BatchedObjects = std::vector<BatchObject>; + static BatchedCommandRequest makeInsertCommand(const NamespaceString& outputNs, + bool bypassDocumentValidation) { + write_ops::InsertCommandRequest insertOp(outputNs); + insertOp.setWriteCommandRequestBase([&] { + write_ops::WriteCommandRequestBase wcb; + wcb.setOrdered(false); + wcb.setBypassDocumentValidation(bypassDocumentValidation); + return wcb; + }()); + return BatchedCommandRequest(std::move(insertOp)); + } + DocumentSourceWriter(const char* stageName, NamespaceString outputNs, const boost::intrusive_ptr<ExpressionContext>& expCtx) : DocumentSource(stageName, expCtx), _outputNs(std::move(outputNs)), - _writeConcern(expCtx->opCtx->getWriteConcern()) {} + _writeConcern(expCtx->opCtx->getWriteConcern()), + _writeSizeEstimator( + expCtx->mongoProcessInterface->getWriteSizeEstimator(expCtx->opCtx, outputNs)) {} DepsTracker::State getDependencies(DepsTracker* deps) const override { deps->needWholeDocument = true; @@ -114,7 +134,7 @@ public: GetModPathsReturn getModifiedPaths() const override { // For purposes of tracking which fields come from where, the writer stage does not modify // any fields by default. - return {GetModPathsReturn::Type::kFiniteSet, std::set<std::string>{}, {}}; + return {GetModPathsReturn::Type::kFiniteSet, OrderedPathSet{}, {}}; } boost::optional<DistributedPlanLogic> distributedPlanLogic() override { @@ -122,7 +142,7 @@ public: } bool canRunInParallelBeforeWriteStage( - const std::set<std::string>& nameOfShardKeyFieldsUponEntryToStage) const override { + const OrderedPathSet& nameOfShardKeyFieldsUponEntryToStage) const override { return true; } @@ -143,9 +163,31 @@ protected: virtual void finalize() {} /** - * Writes the documents in 'batch' to the output namespace. + * Writes the documents in 'batch' to the output namespace via 'bcr'. + */ + virtual void spill(BatchedCommandRequest&& bcr, BatchedObjects&& batch) = 0; + + /** + * Estimates the size of the header of a batch write (that is, the size of the write command + * minus the size of write statements themselves). + */ + int estimateWriteHeaderSize(const BatchedCommandRequest& bcr) const { + using BatchType = BatchedCommandRequest::BatchType; + switch (bcr.getBatchType()) { + case BatchType::BatchType_Insert: + return _writeSizeEstimator->estimateInsertHeaderSize(bcr.getInsertRequest()); + case BatchType::BatchType_Update: + return _writeSizeEstimator->estimateUpdateHeaderSize(bcr.getUpdateRequest()); + case BatchType::BatchType_Delete: + break; + } + MONGO_UNREACHABLE; + } + + /** + * Constructs and configures a BatchedCommandRequest for performing a batch write. */ - virtual void spill(BatchedObjects&& batch) = 0; + virtual BatchedCommandRequest initializeBatchedWriteRequest() const = 0; /** * Creates a batch object from the given document and returns it to the caller along with the @@ -169,6 +211,9 @@ protected: // respect the writeConcern of the original command. WriteConcernOptions _writeConcern; + // An interface that is used to estimate the size of each write operation. + const std::unique_ptr<MongoProcessInterface::WriteSizeEstimator> _writeSizeEstimator; + private: bool _initialized{false}; bool _done{false}; @@ -198,9 +243,30 @@ DocumentSource::GetNextResult DocumentSourceWriter<B>::doGetNext() { _initialized = true; } - BatchedObjects batch; - int bufferedBytes = 0; + // While most metadata attached to a command is limited to less than a KB, Impersonation + // metadata may grow to an arbitrary size. + // + // Ask the active Client how much impersonation metadata we'll use for it, add in our own + // estimate of write header size, and assume that the rest can fit in the space reserved by + // BSONObjMaxUserSize's overhead plus the value from the server parameter: + // internalQueryDocumentSourceWriterBatchExtraReservedBytes. + const auto estimatedMetadataSizeBytes = + rpc::estimateImpersonatedUserMetadataSize(pExpCtx->opCtx); + BatchedCommandRequest batchWrite = initializeBatchedWriteRequest(); + const auto writeHeaderSize = estimateWriteHeaderSize(batchWrite); + const auto initialRequestSize = estimatedMetadataSizeBytes + writeHeaderSize + + internalQueryDocumentSourceWriterBatchExtraReservedBytes.load(); + + uassert(7637800, + "Unable to proceed with write while metadata size ({}KB) exceeds {}KB"_format( + initialRequestSize / 1024, BSONObjMaxUserSize / 1024), + initialRequestSize <= BSONObjMaxUserSize); + + const auto maxBatchSizeBytes = BSONObjMaxUserSize - initialRequestSize; + + BatchedObjects batch; + size_t bufferedBytes = 0; auto nextInput = pSource->getNext(); for (; nextInput.isAdvanced(); nextInput = pSource->getNext()) { waitWhileFailPointEnabled(); @@ -210,16 +276,17 @@ DocumentSource::GetNextResult DocumentSourceWriter<B>::doGetNext() { bufferedBytes += objSize; if (!batch.empty() && - (bufferedBytes > BSONObjMaxUserSize || + (bufferedBytes > maxBatchSizeBytes || batch.size() >= write_ops::kMaxWriteBatchSize)) { - spill(std::move(batch)); + spill(std::move(batchWrite), std::move(batch)); batch.clear(); + batchWrite = initializeBatchedWriteRequest(); bufferedBytes = objSize; } batch.push_back(obj); } if (!batch.empty()) { - spill(std::move(batch)); + spill(std::move(batchWrite), std::move(batch)); batch.clear(); } |
