diff options
Diffstat (limited to 'src/mongo/db/views/resolved_view.cpp')
| -rw-r--r-- | src/mongo/db/views/resolved_view.cpp | 119 |
1 files changed, 80 insertions, 39 deletions
diff --git a/src/mongo/db/views/resolved_view.cpp b/src/mongo/db/views/resolved_view.cpp index 244d38a0b55..6795a874968 100644 --- a/src/mongo/db/views/resolved_view.cpp +++ b/src/mongo/db/views/resolved_view.cpp @@ -33,6 +33,7 @@ #include "mongo/base/init.h" #include "mongo/bson/bsonobjbuilder.h" +#include "mongo/db/pipeline/document_source_coll_stats.h" #include "mongo/db/pipeline/document_source_index_stats.h" #include "mongo/db/pipeline/document_source_internal_convert_bucket_index_stats.h" #include "mongo/db/pipeline/document_source_internal_unpack_bucket.h" @@ -91,11 +92,22 @@ ResolvedView ResolvedView::fromBSON(const BSONObj& commandResponseObj) { mixedSchema = boost::optional<bool>(mixedSchemaElem.boolean()); } + boost::optional<bool> usesExtendedRange = boost::none; + if (auto usesExtendedRangeElem = viewDef[kTimeseriesUsesExtendedRange]) { + uassert(6646910, + str::stream() << "view definition must have " << kTimeseriesUsesExtendedRange + << " of type bool or no such field", + usesExtendedRangeElem.type() == BSONType::Bool); + + usesExtendedRange = boost::optional<bool>(usesExtendedRangeElem.boolean()); + } + return {NamespaceString(viewDef["ns"].valueStringData()), std::move(pipeline), std::move(collationSpec), std::move(timeseriesOptions), - std::move(mixedSchema)}; + std::move(mixedSchema), + std::move(usesExtendedRange)}; } void ResolvedView::serialize(BSONObjBuilder* builder) const { @@ -109,6 +121,10 @@ void ResolvedView::serialize(BSONObjBuilder* builder) const { // Only serialize if it doesn't contain mixed data. if ((_timeseriesMayContainMixedData && !(*_timeseriesMayContainMixedData))) subObj.append(kTimeseriesMayContainMixedData, *_timeseriesMayContainMixedData); + + if ((_timeseriesUsesExtendedRange && (*_timeseriesUsesExtendedRange))) + subObj.append(kTimeseriesUsesExtendedRange, *_timeseriesUsesExtendedRange); + if (!_defaultCollation.isEmpty()) { subObj.append("collation", _defaultCollation); } @@ -118,46 +134,48 @@ std::shared_ptr<const ErrorExtraInfo> ResolvedView::parse(const BSONObj& cmdRepl return std::make_shared<ResolvedView>(fromBSON(cmdReply)); } -AggregateCommandRequest ResolvedView::asExpandedViewAggregation( - const AggregateCommandRequest& request) const { - // Perform the aggregation on the resolved namespace. The new pipeline consists of two parts: - // first, 'pipeline' in this ResolvedView; then, the pipeline in 'request'. - std::vector<BSONObj> resolvedPipeline; - resolvedPipeline.reserve(_pipeline.size() + request.getPipeline().size()); - resolvedPipeline.insert(resolvedPipeline.end(), _pipeline.begin(), _pipeline.end()); - resolvedPipeline.insert( - resolvedPipeline.end(), request.getPipeline().begin(), request.getPipeline().end()); - - // $indexStats needs special handling for time-series-collections. Normally for a regular read, - // $_internalUnpackBucket unpacks the buckets entries into time-series document format and then - // passes the time-series documents on through the pipeline. Instead we need to read the buckets - // collection's index stats unmodified and then pass the results through an additional stage to - // specially convert them to the time-series collection's schema, and then onward. There is no - // need for the $_internalUnpackBucket stage with $indexStats, so we remove it. - if (resolvedPipeline.size() >= 2 && - resolvedPipeline[0][DocumentSourceInternalUnpackBucket::kStageNameInternal] && - resolvedPipeline[1][DocumentSourceIndexStats::kStageName]) { - // Clear the $_internalUnpackBucket stage. - auto unpackStage = resolvedPipeline[0]; - resolvedPipeline[0] = resolvedPipeline[1]; - - // Grab the $_internalUnpackBucket stage's time-series collection schema options and pass - // them into the $_internalConvertBucketIndexStats stage to use for schema conversion. - BSONObjBuilder builder; - for (const auto& elem : - unpackStage[DocumentSourceInternalUnpackBucket::kStageNameInternal].Obj()) { - if (elem.fieldNameStringData() == timeseries::kTimeFieldName || - elem.fieldNameStringData() == timeseries::kMetaFieldName) { +void ResolvedView::handleTimeseriesRewrites(std::vector<BSONObj>* resolvedPipeline) const { + // Stages that are constrained to be the first stage of the pipeline ($collStats, $indexStats) + // require special handling since $_internalUnpackBucket is the first stage. + if (resolvedPipeline->size() >= 2 && + (*resolvedPipeline)[0][DocumentSourceInternalUnpackBucket::kStageNameInternal] && + ((*resolvedPipeline)[1][DocumentSourceIndexStats::kStageName] || + (*resolvedPipeline)[1][DocumentSourceCollStats::kStageName])) { + // Normally for a regular read, $_internalUnpackBucket unpacks the buckets entries into + // time-series document format and then passes the time-series documents on through the + // pipeline. Instead, for $indexStats, we need to read the buckets collection's index + // stats unmodified and then pass the results through an additional stage to specially + // convert them to the time-series collection's schema, and then onward. We grab the + // $_internalUnpackBucket stage's time-series collection schema options and pass them + // into the $_internalConvertBucketIndexStats stage to use for schema conversion. + if ((*resolvedPipeline)[1][DocumentSourceIndexStats::kStageName]) { + auto unpackStage = (*resolvedPipeline)[0]; + (*resolvedPipeline)[0] = (*resolvedPipeline)[1]; + BSONObjBuilder builder; + for (const auto& elem : + unpackStage[DocumentSourceInternalUnpackBucket::kStageNameInternal].Obj()) { + if (elem.fieldNameStringData() == timeseries::kTimeFieldName || + elem.fieldNameStringData() == timeseries::kMetaFieldName) { + builder.append(elem); + } + } + (*resolvedPipeline)[1] = + BSON(DocumentSourceInternalConvertBucketIndexStats::kStageName << builder.obj()); + } else { + auto collStatsStage = (*resolvedPipeline)[1]; + BSONObjBuilder builder; + for (const auto& elem : collStatsStage[DocumentSourceCollStats::kStageName].Obj()) { builder.append(elem); } + builder.append("$_requestOnTimeseriesView", true); + (*resolvedPipeline)[1] = BSON(DocumentSourceCollStats::kStageName << builder.obj()); + // For $collStats, we directly read the collection stats from the buckets + // collection, and skip $_internalUnpackBucket. + resolvedPipeline->erase(resolvedPipeline->begin()); } - resolvedPipeline[1] = - BSON(DocumentSourceInternalConvertBucketIndexStats::kStageName << builder.obj()); - } else if (resolvedPipeline.size() >= 1 && - resolvedPipeline[0][DocumentSourceInternalUnpackBucket::kStageNameInternal] && - serverGlobalParams.featureCompatibility.isGreaterThanOrEqualTo( + } else if (serverGlobalParams.featureCompatibility.isGreaterThanOrEqualTo( multiversion::FeatureCompatibilityVersion::kVersion_5_2)) { - auto unpackStage = resolvedPipeline[0]; + auto unpackStage = (*resolvedPipeline)[0]; BSONObjBuilder builder; for (const auto& elem : @@ -166,11 +184,31 @@ AggregateCommandRequest ResolvedView::asExpandedViewAggregation( } builder.append(DocumentSourceInternalUnpackBucket::kAssumeNoMixedSchemaData, ((_timeseriesMayContainMixedData && !(*_timeseriesMayContainMixedData)))); - resolvedPipeline[0] = + + builder.append(DocumentSourceInternalUnpackBucket::kUsesExtendedRange, + ((_timeseriesUsesExtendedRange && *_timeseriesUsesExtendedRange))); + + (*resolvedPipeline)[0] = BSON(DocumentSourceInternalUnpackBucket::kStageNameInternal << builder.obj()); } +} + +AggregateCommandRequest ResolvedView::asExpandedViewAggregation( + const AggregateCommandRequest& request) const { + // Perform the aggregation on the resolved namespace. The new pipeline consists of two parts: + // first, 'pipeline' in this ResolvedView; then, the pipeline in 'request'. + std::vector<BSONObj> resolvedPipeline; + resolvedPipeline.reserve(_pipeline.size() + request.getPipeline().size()); + resolvedPipeline.insert(resolvedPipeline.end(), _pipeline.begin(), _pipeline.end()); + resolvedPipeline.insert( + resolvedPipeline.end(), request.getPipeline().begin(), request.getPipeline().end()); + + if (resolvedPipeline.size() >= 1 && + resolvedPipeline[0][DocumentSourceInternalUnpackBucket::kStageNameInternal]) { + handleTimeseriesRewrites(&resolvedPipeline); + } - AggregateCommandRequest expandedRequest{_namespace, resolvedPipeline}; + AggregateCommandRequest expandedRequest{_namespace, std::move(resolvedPipeline)}; if (request.getExplain()) { expandedRequest.setExplain(request.getExplain()); @@ -192,6 +230,9 @@ AggregateCommandRequest ResolvedView::asExpandedViewAggregation( } } expandedRequest.setHint(rewritten); + + // JLR Possibly need to handle inappropriately forcing clustered here. Will likely be + // handled by removing the _id predicate, however. } else { expandedRequest.setHint(request.getHint()); } |
