Stages API¶
This page documents all aggregation stage classes.
Document Stages¶
Match
¶
Bases: BaseModel
$match stage - filters documents by specified criteria.
Example
Match(query={"status": "active"}).model_dump() {"$match": {"status": "active"}}
With logical operators¶
Match(query={"$and": [{"status": "active"}, {"age": {"$gt": 18}}]})
Source code in mongo_aggro/stages/core.py
Project
¶
Bases: BaseModel
$project stage - shapes documents by including/excluding fields.
Example
Project(fields={"name": 1, "year": 1, "_id": 0}).model_dump() {"$project": {"name": 1, "year": 1, "_id": 0}}
With expressions¶
Project(fields={"fullName": {"$concat": ["$first", " ", "$last"]}})
Source code in mongo_aggro/stages/core.py
AddFields
¶
Bases: BaseModel
$addFields stage - adds new fields to documents.
Example
AddFields(fields={"isActive": True, "score": {"$sum": "$marks"}}) {"$addFields": {"isActive": true, "score": {"$sum": "$marks"}}}
Source code in mongo_aggro/stages/transform.py
Set
¶
Bases: BaseModel
$set stage - alias for $addFields.
Example
Set(fields={"status": "processed"}).model_dump() {"$set": {"status": "processed"}}
Source code in mongo_aggro/stages/transform.py
Unset
¶
Bases: BaseModel
$unset stage - removes fields from documents.
Example
Unset(fields=["password", "secret"]).model_dump()
Unset(fields="temporaryField").model_dump()
Source code in mongo_aggro/stages/transform.py
ReplaceRoot
¶
Bases: BaseModel
$replaceRoot stage - replaces document with specified embedded document.
Example
ReplaceRoot(new_root="$nested").model_dump() {"$replaceRoot": {"newRoot": "$nested"}}
Source code in mongo_aggro/stages/transform.py
ReplaceWith
¶
Bases: BaseModel
$replaceWith stage - replaces document (alias for $replaceRoot).
Example
ReplaceWith(expression="$embedded").model_dump()
Source code in mongo_aggro/stages/transform.py
Grouping Stages¶
Group
¶
Bases: BaseModel
$group stage - groups documents by specified expression.
Example
Group( ... id="$category", ... total={"$sum": "$quantity"}, ... count={"$sum": 1} ... ).model_dump() { "$group": { "_id": "$category", "total": {"$sum": "$quantity"}, "count": {"$sum": 1} } }
Source code in mongo_aggro/stages/core.py
SortByCount
¶
Bases: BaseModel
$sortByCount stage - groups and counts by field, sorted by count.
Example
SortByCount(field="category").model_dump()
Source code in mongo_aggro/stages/group.py
Bucket
¶
Bases: BaseModel
$bucket stage - categorizes documents into buckets.
Example
Bucket( ... group_by="$price", ... boundaries=[0, 100, 500, 1000], ... default="Other", ... output={"count": {"$sum": 1}} ... ).model_dump() {"$bucket": { "groupBy": "$price", "boundaries": [0, 100, 500, 1000], "default": "Other", "output": {"count": {"$sum": 1}} }}
Source code in mongo_aggro/stages/group.py
BucketAuto
¶
Bases: BaseModel
$bucketAuto stage - automatically categorizes into specified buckets.
Example
BucketAuto(group_by="$age", buckets=5).model_dump() {"$bucketAuto": {"groupBy": "$age", "buckets": 5}}
Source code in mongo_aggro/stages/group.py
Array Stages¶
Unwind
¶
Bases: BaseModel
$unwind stage - deconstructs an array field.
Example
Unwind(path="cars").model_dump()
With options¶
Unwind( ... path="items", ... include_array_index="itemIndex", ... preserve_null_and_empty=True ... ).model_dump() {"$unwind": { "path": "$items", "includeArrayIndex": "itemIndex", "preserveNullAndEmptyArrays": true }}
Source code in mongo_aggro/stages/array.py
Join Stages¶
Lookup
¶
Bases: BaseModel
$lookup stage - performs a left outer join.
Example
Simple lookup¶
Lookup( ... from_collection="products", ... local_field="product_id", ... foreign_field="_id", ... as_field="product" ... ).model_dump() {"$lookup": { "from": "products", "localField": "product_id", "foreignField": "_id", "as": "product" }}
With pipeline¶
Lookup( ... from_collection="orders", ... let={"customerId": "$_id"}, ... pipeline=Pipeline([Match(query={"status": "active"})]), ... as_field="orders" ... ).model_dump()
Source code in mongo_aggro/stages/join.py
GraphLookup
¶
Bases: BaseModel
$graphLookup stage - performs recursive search.
Example
GraphLookup( ... from_collection="employees", ... start_with="$reportsTo", ... connect_from_field="reportsTo", ... connect_to_field="name", ... as_field="reportingHierarchy" ... ).model_dump()
Source code in mongo_aggro/stages/join.py
135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 | |
Sorting and Pagination¶
Sort
¶
Bases: BaseModel
$sort stage - sorts documents.
Example
Sort(fields={"age": -1, "name": 1}).model_dump() {"$sort": {"age": -1, "name": 1}}
Source code in mongo_aggro/stages/core.py
Limit
¶
Bases: BaseModel
$limit stage - limits the number of documents.
Example
Limit(count=10).model_dump()
Source code in mongo_aggro/stages/core.py
Skip
¶
Bases: BaseModel
$skip stage - skips a number of documents.
Example
Skip(count=5).model_dump()
Source code in mongo_aggro/stages/core.py
Sample
¶
Bases: BaseModel
$sample stage - randomly selects documents.
Example
Sample(size=10).model_dump() {"$sample": {"size": 10}}
Source code in mongo_aggro/stages/output.py
Multi-Pipeline Stages¶
Facet
¶
Bases: BaseModel
$facet stage - processes multiple pipelines within a single stage.
Example
Facet(pipelines={ ... "byCategory": Pipeline([Group(id="$category")]), ... "byYear": Pipeline([Group(id="$year")]) ... }).model_dump() {"$facet": { "byCategory": [{"$group": {"_id": "$category"}}], "byYear": [{"$group": {"_id": "$year"}}] }}
Source code in mongo_aggro/stages/group.py
UnionWith
¶
Bases: BaseModel
$unionWith stage - combines pipeline results with another collection.
Example
UnionWith(collection="archive").model_dump()
UnionWith( ... collection="archive", ... pipeline=Pipeline([Match(query={"year": 2023})]) ... ).model_dump() {"$unionWith": {"coll": "archive", "pipeline": [...]}}
Source code in mongo_aggro/stages/join.py
Output Stages¶
Count
¶
Bases: BaseModel
$count stage - counts documents.
Example
Count(field="total").model_dump()
Source code in mongo_aggro/stages/core.py
Out
¶
Bases: BaseModel
$out stage - writes results to a collection.
Example
Out(collection="results").model_dump()
Out(collection="results", db="analytics").model_dump() {"$out": {"db": "analytics", "coll": "results"}}
Source code in mongo_aggro/stages/output.py
Merge
¶
Bases: BaseModel
$merge stage - writes results to a collection with merge behavior.
Example
Merge( ... into="reports", ... on="_id", ... when_matched="merge", ... when_not_matched="insert" ... ).model_dump() {"$merge": { "into": "reports", "on": "_id", "whenMatched": "merge", "whenNotMatched": "insert" }}
Source code in mongo_aggro/stages/output.py
Other Stages¶
Redact
¶
Bases: BaseModel
$redact stage - restricts document content based on stored info.
Example
Redact(expression={ ... "$cond": { ... "if": {"$eq": ["$level", 5]}, ... "then": "$$PRUNE", ... "else": "$$DESCEND" ... } ... }).model_dump()
Source code in mongo_aggro/stages/transform.py
GeoNear
¶
Bases: BaseModel
$geoNear stage - returns documents near a geographic point.
Example
GeoNear( ... near={"type": "Point", "coordinates": [-73.99, 40.73]}, ... distance_field="dist.calculated", ... spherical=True, ... max_distance=5000 ... ).model_dump()
Source code in mongo_aggro/stages/geo.py
SetWindowFields
¶
Bases: BaseModel
$setWindowFields stage - performs window calculations.
Example
SetWindowFields( ... partition_by="$state", ... sort_by={"date": 1}, ... output={ ... "cumulative": { ... "$sum": "$quantity", ... "window": {"documents": ["unbounded", "current"]} ... } ... } ... ).model_dump()
Source code in mongo_aggro/stages/window.py
Densify
¶
Bases: BaseModel
$densify stage - fills gaps in data.
Example
Densify( ... field="date", ... range={"step": 1, "unit": "day", "bounds": "full"}, ... partition_by_fields=["series"] ... ).model_dump()
Source code in mongo_aggro/stages/window.py
Fill
¶
Bases: BaseModel
$fill stage - fills null/missing field values.
Example
Fill( ... sort_by={"date": 1}, ... output={ ... "score": {"method": "linear"}, ... "bootcamp": {"value": "missing"} ... } ... ).model_dump()
Source code in mongo_aggro/stages/window.py
Documents
¶
Bases: BaseModel
$documents stage - returns literal documents.
Example
Documents(documents=[ ... {"x": 1, "y": 2}, ... {"x": 3, "y": 4} ... ]).model_dump() {"$documents": [{"x": 1, "y": 2}, {"x": 3, "y": 4}]}
Source code in mongo_aggro/stages/output.py
Statistics Stages¶
CollStats
¶
Bases: BaseModel
$collStats stage - returns collection statistics.
Example
CollStats(lat_stats={"histograms": True}).model_dump() {"$collStats": {"latencyStats": {"histograms": True}}}
CollStats(storage_stats={}).model_dump() {"$collStats": {"storageStats": {}}}
CollStats(count={}).model_dump() {"$collStats": {"count": {}}}
Source code in mongo_aggro/stages/stats.py
IndexStats
¶
Bases: BaseModel
$indexStats stage - returns index usage statistics.
Example
IndexStats().model_dump() {"$indexStats": {}}
Source code in mongo_aggro/stages/stats.py
PlanCacheStats
¶
Bases: BaseModel
$planCacheStats stage - returns plan cache information.
Example
PlanCacheStats().model_dump() {"$planCacheStats": {}}
Source code in mongo_aggro/stages/stats.py
Session Stages¶
ListSessions
¶
Bases: BaseModel
$listSessions stage - lists all sessions in system.sessions.
Example
ListSessions().model_dump() {"$listSessions": {}}
ListSessions(users=[{"user": "admin", "db": "admin"}]).model_dump() {"$listSessions": {"users": [{"user": "admin", "db": "admin"}]}}
ListSessions(all_users=True).model_dump() {"$listSessions": {"allUsers": True}}
Source code in mongo_aggro/stages/session.py
ListLocalSessions
¶
Bases: BaseModel
$listLocalSessions stage - lists local sessions (db.aggregate only).
Example
ListLocalSessions().model_dump() {"$listLocalSessions": {}}
ListLocalSessions(all_users=True).model_dump() {"$listLocalSessions": {"allUsers": True}}
Source code in mongo_aggro/stages/session.py
ListSampledQueries
¶
Bases: BaseModel
$listSampledQueries stage - lists sampled queries.
Example
ListSampledQueries().model_dump() {"$listSampledQueries": {}}
ListSampledQueries(namespace="db.collection").model_dump() {"$listSampledQueries": {"namespace": "db.collection"}}
Source code in mongo_aggro/stages/session.py
Change Stream Stages¶
ChangeStream
¶
Bases: BaseModel
$changeStream stage - returns a change stream cursor.
Must be the first stage in the pipeline.
Example
ChangeStream().model_dump() {"$changeStream": {}}
ChangeStream(full_document="updateLookup").model_dump() {"$changeStream": {"fullDocument": "updateLookup"}}
Source code in mongo_aggro/stages/change.py
ChangeStreamSplitLargeEvent
¶
Bases: BaseModel
$changeStreamSplitLargeEvent stage - splits large change events.
Must be the last stage in a $changeStream pipeline.
Example
ChangeStreamSplitLargeEvent().model_dump() {"$changeStreamSplitLargeEvent": {}}
Source code in mongo_aggro/stages/change.py
Admin Stages¶
CurrentOp
¶
Bases: BaseModel
$currentOp stage - returns current operations (db.aggregate only).
Example
CurrentOp().model_dump() {"$currentOp": {}}
CurrentOp(all_users=True, idle_connections=True).model_dump() {"$currentOp": {"allUsers": True, "idleConnections": True}}
Source code in mongo_aggro/stages/stats.py
ListClusterCatalog
¶
Bases: BaseModel
$listClusterCatalog stage - lists collections in a cluster.
Example
ListClusterCatalog().model_dump() {"$listClusterCatalog": {}}
Source code in mongo_aggro/stages/misc.py
ListSearchIndexes
¶
Bases: BaseModel
$listSearchIndexes stage - lists Atlas Search indexes.
Example
ListSearchIndexes().model_dump() {"$listSearchIndexes": {}}
ListSearchIndexes(id="index_id").model_dump() {"$listSearchIndexes": {"id": "index_id"}}
ListSearchIndexes(name="index_name").model_dump() {"$listSearchIndexes": {"name": "index_name"}}
Source code in mongo_aggro/stages/search.py
Atlas Search Stages¶
Search
¶
Bases: BaseModel
$search stage - Atlas full-text search.
Example
Search(index="default", text={"query": "coffee", "path": "title"}) {"$search": {"index": "default", "text": {"query": "coffee", ...}}}
Source code in mongo_aggro/stages/search.py
47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 | |
SearchMeta
¶
Bases: BaseModel
$searchMeta stage - returns Atlas Search metadata.
Example
SearchMeta(index="default", count={"type": "total"}).model_dump() {"$searchMeta": {"index": "default", "count": {"type": "total"}}}
Source code in mongo_aggro/stages/search.py
VectorSearch
¶
Bases: BaseModel
$vectorSearch stage - Atlas vector search (MongoDB 7.0.2+).
Example
VectorSearch( ... index="vector_index", ... path="embedding", ... query_vector=[0.1, 0.2, 0.3], ... num_candidates=100, ... limit=10 ... ).model_dump() {"$vectorSearch": { "index": "vector_index", "path": "embedding", "queryVector": [0.1, 0.2, 0.3], "numCandidates": 100, "limit": 10 }}
Source code in mongo_aggro/stages/search.py
Advanced Stages¶
QuerySettings
¶
Bases: BaseModel
$querySettings stage - returns query settings (MongoDB 8.0+).
Example
QuerySettings().model_dump() {"$querySettings": {}}
Source code in mongo_aggro/stages/misc.py
RankFusion
¶
Bases: BaseModel
$rankFusion stage - combines ranked results from multiple pipelines.
Example
RankFusion( ... input={"search": [...], "vector": [...]}, ... combination={"weights": {"search": 0.7, "vector": 0.3}} ... ).model_dump() {"$rankFusion": { "input": {"search": [...], "vector": [...]}, "combination": {"weights": {"search": 0.7, "vector": 0.3}} }}