Skip to content

Aggregation Pipeline

Build MongoDB aggregation pipelines with a fluent builder.

Setup

from vellum import AggregationPipeline

pipeline = AggregationPipeline(repo.collection)

Stages

$match

results = await pipeline.match({"status": "active"}).execute()

$group

results = await (
    pipeline
    .group("$category", total={"$sum": 1}, avg_price={"$avg": "$price"})
    .execute()
)

$project

results = await pipeline.project({"name": 1, "price": 1, "_id": 0}).execute()

$sort

from vellum import FieldRef

quantity = FieldRef("quantity")
results = await pipeline.sort(quantity.desc()).execute()
# Or: pipeline.sort([("quantity", -1)])

$limit / $skip

results = await pipeline.sort(quantity.desc()).limit(5).skip(10).execute()

$unwind

results = await pipeline.unwind("$items").execute()
# With options: pipeline.unwind({"path": "$items", "preserveNullAndEmptyArrays": True})

$lookup (join)

results = await (
    pipeline
    .lookup(
        from_collection="orders",
        local_field="user_id",
        foreign_field="_id",
        as_field="orders",
    )
    .execute()
)

$addFields

results = await pipeline.add_fields({"total": {"$multiply": ["$price", "$quantity"]}}).execute()

$count

results = await pipeline.count("total").execute()
# Returns [{"total": 42}]

Output Models

Validate aggregation results with Pydantic models:

from pydantic import BaseModel

class RevenueByCategory(BaseModel):
    category: str
    total_revenue: float

results = await (
    pipeline
    .group("$category", total_revenue={"$sum": "$revenue"})
    .sort([("total_revenue", -1)])
    .execute(output_model=RevenueByCategory)
)
# results: list[RevenueByCategory]