---
title: Aggregation Pipeline
description: Run a MongoDB-compatible subset of aggregation stages, group accumulators, and a recursive expression engine over an OMGDB collection.
---

An aggregation pipeline transforms a stream of documents through an ordered list of stages. In OMGDB a pipeline is a JSON **array** of stage documents — each stage is an object with exactly one `$`-prefixed key — passed to the CLI as a single JSON-string argument:

```sh
omgdb aggregate app.omgdb <collection> '[<stage>, <stage>, ...]'
```

The pipeline runs against the documents scanned from `<collection>`, then each stage feeds its output documents into the next. Stages, group accumulators, and expression operators are each dispatched by a single flat `match` on the operator name in the `omgdb-agg` crate. The design is deliberately *additive* — supporting a new operator means adding one match arm — but it is not a runtime registry: there is no registration API, trait objects, or plugin mechanism, and an unrecognized operator returns an unknown-operator error (with a did-you-mean suggestion when a close match exists).

Expressions are evaluated by a recursive engine. A string beginning with `$` (e.g. `"$a.b"`) is a **field path** resolved against the current document (missing paths yield `null`); a string beginning with `$$` is a [system variable](#system-variables); any other string is a literal. An object whose first key starts with `$` is a single-operator expression (and must have exactly one key); any other object is a literal sub-document evaluated field by field; arrays are evaluated element-wise.

This is a **compatible subset** of MongoDB's aggregation framework, not a complete implementation. See [query operators](/docs/query-operators/) for the filter syntax used inside `$match`.

## Stages

Stages run in array order. Each stage object must have exactly one `$`-key; otherwise the pipeline is malformed.

| Stage | Description | Example |
| --- | --- | --- |
| `$documents` | Injects an inline array of documents as the pipeline input. Must be the **first** stage. | `{"$documents":[{"a":1},{"a":2}]}` |
| `$match` | Keeps documents matching a query filter (compiled via the query engine). May include an `$expr` expression predicate alongside ordinary field conditions. | `{"$match":{"age":{"$gte":25}}}` |
| `$project` | Inclusion/exclusion projection plus computed fields. `1`/`true` includes, `0`/`false` excludes, any other value is an expression. | `{"$project":{"_id":0,"full":1}}` |
| `$addFields` / `$set` | Adds or overwrites top-level fields from evaluated expressions. `$set` is an exact alias. | `{"$addFields":{"sum":{"$add":["$a","$b"]}}}` |
| `$unset` | Removes one field (string) or several (array of strings); dotted paths remove nested fields. | `{"$unset":["tmp","meta.debug"]}` |
| `$group` | Groups by an `_id` expression and computes accumulators per group. | `{"$group":{"_id":"$dept","total":{"$sum":"$sal"}}}` |
| `$bucket` | Buckets documents by a `groupBy` expression into explicit `boundaries`, with optional `default` and `output` accumulators. | `{"$bucket":{"groupBy":"$age","boundaries":[0,18,65,120]}}` |
| `$sortByCount` | Groups by an expression and emits `{_id, count}` documents, most frequent first. | `{"$sortByCount":"$dept"}` |
| `$sort` | Stable-sorts by one or more keys; each direction must be `1` or `-1`. | `{"$sort":{"age":-1}}` |
| `$limit` | Keeps at most the first N documents. | `{"$limit":1}` |
| `$skip` | Drops the first N documents. | `{"$skip":1}` |
| `$count` | Collapses the stream to one document with the named field set to the document count. | `{"$count":"n"}` |
| `$unwind` | Emits one document per element of an array field. | `{"$unwind":"$tags"}` |
| `$replaceRoot` / `$replaceWith` | Promotes an embedded document (or an expression that evaluates to an object) to the root. | `{"$replaceRoot":{"newRoot":"$meta"}}` |
| `$lookup` | Left-outer join of another collection: equality (`localField`/`foreignField`), uncorrelated `pipeline`, or correlated `let` + `pipeline` with `$expr`. Requires a store. | `{"$lookup":{"from":"items","localField":"item","foreignField":"_id","as":"itemDocs"}}` |
| `$unionWith` | Appends all documents from another collection (optionally after a sub-pipeline). Duplicates are preserved — UNION ALL semantics. | `{"$unionWith":"archive"}` |
| `$facet` | Runs several named sub-pipelines over the same input and emits one document of result arrays. | `{"$facet":{"count":[{"$count":"n"}]}}` |

### Stage notes

- **`$documents`** requires an array whose entries are all objects, and it is rejected anywhere except position 0.
- **`$project`** operates on top-level fields. An inclusion projection keeps `_id` first unless `_id:0`; in an inclusion projection, excluded fields other than `_id` are simply ignored rather than raising an error. A projection with only `0`/`false` values is an exclusion projection.
- **`$unset`** takes a field-name string or an array of them. Names may be dotted paths (`meta.debug`); removing through an array removes the key from every element. Names must be nonempty and must not start with `$`.
- **`$group`** requires an `_id`. Every non-`_id` field must be a single-accumulator object. Groups are keyed by the canonical JSON of the `_id` value and stored in an ordered map, so output is ordered by that canonical key rather than by input order.
- **`$bucket`** requires `groupBy` and a `boundaries` array of at least two values that share one type (mixed integer/float is allowed) and are strictly ascending. A document whose `groupBy` value falls outside every boundary goes to the `default` bucket — the `default` id must itself sort outside the boundary range — or, with no `default`, is a type error. `output` declares accumulators per bucket (the default output is `{"count":{"$sum":1}}`); empty buckets are omitted.
- **`$sortByCount`** groups by any expression (a document literal is rejected), counts each group with checked arithmetic, and sorts by count descending with ties broken by `_id` under the engine's value order.
- **`$sort`** keys support dotted paths. Missing field values sort before present ones. A direction other than `1` or `-1` is a malformed error.
- **`$unwind`** accepts a `"$field"` string or `{"path":"$field"}`. It operates on a top-level field (the leading `$` is stripped). A missing or `null` array value silently drops the document; a non-array value is passed through as a single document. There are no `preserveNullAndEmptyArrays` or `includeArrayIndex` options.
- **`$replaceRoot`** errors if `newRoot` does not evaluate to an object. `$replaceWith: <expr>` is sugar for `$replaceRoot: { newRoot: <expr> }`.
- **`$unionWith`** accepts a collection-name string, or `{coll, pipeline}` where the optional `pipeline` runs over the unioned collection first. `coll` may be omitted only when the pipeline starts with `$documents` (an inline union).
- **`$facet`** collapses the stream to a single document. Sub-pipelines inherit the store, so they may use `$lookup` and `$unionWith`.

### $lookup in depth

`$lookup` supports all three MongoDB join forms:

```json
{"$lookup":{"from":"items","localField":"item","foreignField":"_id","as":"itemDocs"}}
{"$lookup":{"from":"items","pipeline":[{"$match":{"active":true}}],"as":"activeItems"}}
{"$lookup":{"from":"orders","let":{"uid":"$_id"},
  "pipeline":[{"$match":{"$expr":{"$eq":["$userId","$$uid"]}}}],"as":"orders"}}
```

- The **equality form** builds a hash map over the foreign side (one pass, not a nested scan per input document) and is array-aware on both sides: a scalar on either side matches an element of an array on the other. Dotted `localField`/`foreignField` paths are supported.
- The **uncorrelated pipeline form** runs `pipeline` over the foreign collection once and attaches the result to every input document.
- The **correlated form** binds `let` variables per input document; the sub-pipeline reads them as `$$name` (typically inside `$match` + `$expr`). Variable names must be nonempty and must not start with `$` or contain `.`.
- `localField` and `foreignField` must appear together, and at least one of `localField`/`pipeline` is required. Unknown options are rejected by name.

> **Note:** `$lookup` and a collection-sourced `$unionWith` join other collections and therefore need a backing store. They are available when the pipeline runs against a collection (as the `omgdb aggregate` command does); running a pipeline over an in-memory document list without a store makes them malformed errors.

## Accumulators

Accumulators appear inside `$group` (and `$bucket`'s `output`), one per output field, applied to the values produced by evaluating the accumulator's argument expression for each document in the group.

| Accumulator | Description | Example |
| --- | --- | --- |
| `$sum` | Sums numeric inputs; all-integer inputs yield an integer (checked, overflow is an error), otherwise a float. Non-numeric values are ignored. Also used as a counter via `{"$sum":1}`. | `{"total":{"$sum":"$sal"}}` |
| `$avg` | Mean of numeric inputs as a float; ignores non-numeric values; `null` if there are none. | `{"avg":{"$avg":"$sal"}}` |
| `$min` | Smallest value by the engine's comparison ordering; `null` for empty input. | `{"lo":{"$min":"$sal"}}` |
| `$max` | Largest value by the comparison ordering; `null` for empty input. | `{"hi":{"$max":"$sal"}}` |
| `$first` | First accumulated value in the group (input order); `null` if empty. | `{"f":{"$first":"$sal"}}` |
| `$last` | Last accumulated value in the group; `null` if empty. | `{"l":{"$last":"$sal"}}` |
| `$push` | Collects all accumulated values into an array. | `{"all":{"$push":"$sal"}}` |
| `$addToSet` | Collects distinct values into an array (deduplicated by canonical value, first-seen order preserved). | `{"tags":{"$addToSet":"$tag"}}` |
| `$mergeObjects` | Merges accumulated object values into one document; later fields overwrite earlier ones. `null` inputs are skipped; a non-object input is a type error. | `{"merged":{"$mergeObjects":"$meta"}}` |
| `$stdDevPop` | Population standard deviation of numeric inputs as a float; `null` for empty input. | `{"sd":{"$stdDevPop":"$sal"}}` |
| `$stdDevSamp` | Sample standard deviation; `null` with fewer than two numeric inputs. | `{"sd":{"$stdDevSamp":"$sal"}}` |

## Expression operators

Expression operators are used inside `$project`, `$addFields`/`$set`, `$group` accumulator arguments, `$replaceRoot`, `$bucket`'s `groupBy`, `$sortByCount`, and `$match`'s `$expr`. They take a single argument or an argument array, each element of which is itself an expression evaluated against the current document.

### System variables

A string beginning with `$$` reads a system variable. `$$ROOT` and `$$CURRENT` both resolve to the current document, and both accept dotted paths (`"$$ROOT.a.b"`). Inside a correlated `$lookup` sub-pipeline, the join's `let` bindings are also available as `$$name` (with optional dotted paths into object values). Any other `$$`-name is a malformed error.

### Literal and conditional

| Operator | Description | Example |
| --- | --- | --- |
| `$literal` | Returns its operand unevaluated, so `$`-prefixed strings and operator objects are treated as literal data. | `{"$literal":"$notAFieldPath"}` |
| `$cond` | Ternary: returns the `then` branch when the condition is truthy, else `else`. Accepts `[if, then, else]` or `{if, then, else}`. | `{"$cond":[{"$gte":["$age",18]},true,false]}` |
| `$ifNull` | Returns the first argument unless it is `null`, in which case the second. Exactly 2 arguments. | `{"$ifNull":["$nickname","anon"]}` |
| `$switch` | Evaluates `branches` in order, returning the `then` of the first truthy `case`; falls back to `default`. | `{"$switch":{"branches":[{"case":{"$gte":["$score",90]},"then":"A"}],"default":"F"}}` |
| `$type` | BSON-style type name of the argument as a string (e.g. an integer reports `"long"`). | `{"$type":"$score"}` |

> **Note:** No matching `$switch` branch and no `default` is a type error.

### Arithmetic

All arithmetic operators require numeric arguments; a non-numeric argument is a type error.

| Operator | Description | Example |
| --- | --- | --- |
| `$add` | Adds arguments; all-integer args stay integer (checked overflow), otherwise float. | `{"$add":["$a","$b"]}` |
| `$subtract` | First minus second; integer when both integral (checked), else float. Exactly 2 arguments. | `{"$subtract":["$a","$b"]}` |
| `$multiply` | Multiplies arguments; integer when all integral (checked overflow), else float. | `{"$multiply":[2,3]}` |
| `$divide` | First divided by second; always returns a float. Divide-by-zero is a type error. Exactly 2 arguments. | `{"$divide":["$a",2]}` |
| `$mod` | Remainder of first by second; integer when both integral (checked), else float. Mod-by-zero is a type error. | `{"$mod":["$a",-1]}` |

### Comparison

Each comparison takes exactly 2 arguments and returns a boolean, comparing under the engine's ordering.

| Operator | Description | Example |
| --- | --- | --- |
| `$eq` | True when the arguments compare equal. | `{"$eq":["$a",1]}` |
| `$ne` | True when the arguments are not equal. | `{"$ne":["$a",1]}` |
| `$gt` | True when the first is greater than the second. | `{"$gt":["$score",90]}` |
| `$gte` | True when the first is greater than or equal to the second. | `{"$gte":["$score",90]}` |
| `$lt` | True when the first is less than the second. | `{"$lt":["$a",10]}` |
| `$lte` | True when the first is less than or equal to the second. | `{"$lte":["$a",10]}` |

### Boolean logic

A value is *falsy* if it is `null`, `false`, integer `0`, or float `0.0`; everything else (including empty string or array) is truthy.

| Operator | Description | Example |
| --- | --- | --- |
| `$and` | True when all arguments are truthy. | `{"$and":[{"$gte":["$a",1]},{"$lt":["$a",9]}]}` |
| `$or` | True when any argument is truthy. | `{"$or":[{"$eq":["$a",1]},{"$eq":["$a",2]}]}` |
| `$not` | Negates the truthiness of its (first) argument. | `{"$not":["$flag"]}` |

> **Note:** `$and` and `$or` do **not** short-circuit — all arguments are evaluated before the result is combined.

### Math

| Operator | Description | Example |
| --- | --- | --- |
| `$abs` | Absolute value; integer stays integer (checked overflow), float stays float. | `{"$abs":-7}` |
| `$ceil` | Smallest integer-valued result not less than the number; an integer is returned unchanged. | `{"$ceil":2.1}` |
| `$floor` | Largest integer-valued result not greater than the number; an integer is returned unchanged. | `{"$floor":2.9}` |
| `$round` | Rounds to the nearest integer value (ties away from zero); an integer is returned unchanged. | `{"$round":2.5}` |
| `$trunc` | Truncates toward zero to an integer value; an integer is returned unchanged. | `{"$trunc":2.9}` |
| `$sqrt` | Square root as a float. | `{"$sqrt":16}` |

> **Limitation:** `$round` and `$trunc` take no precision/decimal-places argument (unlike MongoDB). `$ceil`/`$floor`/`$round`/`$trunc` applied to a float return a float with an integral value (e.g. `3.0`), not an integer-typed value.

### Strings

| Operator | Description | Example |
| --- | --- | --- |
| `$concat` | Concatenates string arguments; returns `null` if any argument is `null`; a non-string, non-null argument is a type error. | `{"$concat":["$first"," ","$last"]}` |
| `$toUpper` | Uppercases a string; `null`/missing yields an empty string. | `{"$toUpper":"$name"}` |
| `$toLower` | Lowercases a string; `null`/missing yields an empty string. | `{"$toLower":"$name"}` |
| `$strLenCP` | Number of Unicode code points in a string; a non-string argument is a type error. | `{"$strLenCP":"héllo"}` |
| `$split` | Splits a string by a delimiter into an array; an empty delimiter returns the whole string as a single element. Requires `[string, string]`. | `{"$split":["a,b,c",","]}` |

### Arrays and objects

| Operator | Description | Example |
| --- | --- | --- |
| `$size` | Length of an array as an integer; a non-array argument is a type error. | `{"$size":"$tags"}` |
| `$arrayElemAt` | Element at an integer index; negative indexes count from the end; out-of-range returns `null`. Requires `[array, integer]`. | `{"$arrayElemAt":["$tags",-1]}` |
| `$in` | True when the first argument equals any element of the second (array) argument. Requires `[value, array]`. | `{"$in":["b",["a","b"]]}` |
| `$isArray` | True when the argument is an array. | `{"$isArray":"$tags"}` |
| `$concatArrays` | Concatenates array arguments into one array; a non-array argument is a type error. | `{"$concatArrays":[[1,2],[3]]}` |
| `$mergeObjects` | Merges object arguments into one document; later fields overwrite earlier ones. `null` arguments are skipped; anything else is a type error. | `{"$mergeObjects":["$defaults","$overrides"]}` |

> **Limitation:** Genuinely missing from the expression engine today: `$map`, `$filter`, `$reduce`, `$arrayToObject`, `$reverseArray`, `$sortArray`, the date operators, `$regexMatch`, `$substr`/`$substrCP`, `$trim`, and the type-conversion operators (`$toString`, `$toInt`, `$convert`, …). The `$$NOW` system variable is also unsupported. Each returns an unknown-operator or malformed error rather than a wrong answer.

> **Limitation:** Missing stages include `$bucketAuto`, `$out`, `$merge`, `$sample`, and `$graphLookup` — an unrecognized stage returns an unknown-operator error (with a did-you-mean hint when a close match exists).

## Worked example

Group by a field with `$sum` and `$avg`, then sort the groups. Given a `salaries` collection of `{"dept": ..., "sal": ...}` documents:

```sh
omgdb aggregate app.omgdb salaries '[
  {"$group":{"_id":"$dept","total":{"$sum":"$sal"},"avg":{"$avg":"$sal"},"n":{"$sum":1}}},
  {"$sort":{"total":-1}}
]'
```

For input documents:

```json
[
  {"dept":"a","sal":100},
  {"dept":"a","sal":200},
  {"dept":"b","sal":50}
]
```

The `$group` stage produces one document per department, then `$sort` orders them by `total` descending:

```json
[
  {"_id":"a","total":300,"avg":150.0,"n":2},
  {"_id":"b","total":50,"avg":50.0,"n":1}
]
```

Note that `total` is an integer (`300`) because every summed value was an integer, while `avg` is a float (`150.0`). The `{"$sum":1}` accumulator counts documents per group.

A second example combines computed fields with a projection — adding a concatenated `full` field, then keeping only it:

```sh
omgdb aggregate app.omgdb people '[
  {"$addFields":{"full":{"$concat":["$first"," ","$last"]}}},
  {"$project":{"_id":0,"full":1}}
]'
```

For `{"first":"ada","last":"lovelace"}` this yields `{"full":"ada lovelace"}`.

## Error behavior

Pipeline errors are surfaced rather than silently swallowed:

- A pipeline that is not an array, or a stage that is not a single `$`-operator object, is a **malformed** error.
- An unrecognized stage, accumulator, or expression operator is an **unknown operator** error, and when a known operator is within edit distance 2 the message adds a did-you-mean suggestion — so a typo like `$grup` comes back as ``$grup (stage) (did you mean `$group`?)`` and an agent can repair the pipeline in one step.
- An operand of the wrong type (e.g. `$concat` on a number, `$divide` by zero, integer overflow in `$add`/`$subtract`/`$multiply`/`$mod`/`$sum`) is a **type** error. Integer arithmetic uses checked operations, so overflow returns an error instead of panicking.
- A `$match` filter that fails to compile surfaces the underlying query error.
