skip to content

What input does an Elasticsearch pipeline aggregation consume, and when does it run?

level: juniorimportance: should knowfreq 42%

answer

  1. Not documents — something else's output
  2. Pointed at its input by name
  3. Runs after the other aggregation finishes
  4. Coordinating node, reduce phase
  5. Two families: inside vs beside

basics

~20 s

A pipeline aggregation consumes the output of another aggregation — its buckets or its metric values — instead of documents. It is pointed at that output with buckets_path and is computed during the reduce phase, after the source aggregation has produced results.

solid answer

~40 s

Metric and bucket aggregations read documents; a pipeline aggregation reads the **output** of another aggregation. You point it at that output with `buckets_path`. It then either adds a value inside each bucket of the aggregation it is nested in (a parent pipeline such as `derivative`, `cumulative_sum` or `bucket_script`) or produces one new result beside the aggregation it summarises (a sibling pipeline such as `avg_bucket` or `max_bucket`). Because it needs finished results, Elasticsearch computes it in the reduce phase on the coordinating node, after per-shard aggregation results are merged — it never touches documents and never causes extra document collection. Pipeline aggregations accept no sub-aggregations, but they can chain: one may read another pipeline aggregation's output.

go deeper

for a junior

Be ready to say in one sentence that a pipeline aggregation takes another aggregation's output as input, and to name one example such as derivative or avg_bucket.

for a middle

Explain the mechanics: buckets_path names the input, the aggregation runs in the reduce phase after shard results merge, and it emits either a per-bucket value or one summary value.

for a senior

Show you know the operational consequence — pipeline aggregations cost coordinating-node work, never reduce shard work, and can only see the buckets the source aggregation actually returned.

for a principal

Own the boundary question: which derived numbers belong in a query-time pipeline aggregation versus being precomputed at index time, given bucket-count growth and coordinating-node pressure on shared clusters.

## The three kinds of aggregation Elasticsearch aggregations come in three families. **Metric** aggregations (`avg`, `sum`, `min`, `max`, `cardinality`) turn the documents in a bucket into a number. **Bucket** aggregations (`terms`, `date_histogram`, `range`, `filters`) group matching documents into buckets, each with a `doc_count`, and may contain sub-aggregations. Both of these read documents. A **pipeline** aggregation reads neither documents nor field values. Its input is the *result* of another aggregation that has already run. That single fact is the whole definition, and every other property follows from it. ## buckets_path: naming the input Every pipeline aggregation carries a `buckets_path` that names what it consumes. In its simplest form it is the name of a metric aggregation living in the same bucket: ```json "sales_change": { "derivative": { "buckets_path": "sales" } } ``` The path language also has a `>` separator for descending into sub-aggregations (`sales_per_month>sales`), a `.` for picking one value out of a multi-value metric (`sales_stats.sum`), and special keywords such as `_count` for a bucket's document count and `_key` for its key. Aggregations that combine several inputs (`bucket_script`, `bucket_selector`) take `buckets_path` as a map of names to paths and expose them to a Painless script as `params`. ## The two families **Parent pipelines** are declared *inside* the multi-bucket aggregation whose buckets they process, and they act on each of those buckets: `derivative` and `serial_diff` (change between consecutive buckets), `cumulative_sum` (running total), `moving_fn` (a windowed function such as a moving average), `bucket_script` (a computed value per bucket), `bucket_selector` (drops buckets), `bucket_sort` (reorders and truncates buckets). **Sibling pipelines** are declared *beside* the aggregation they consume, at the same level, and produce a single new result over all of its buckets: `avg_bucket`, `sum_bucket`, `min_bucket`, `max_bucket`, `stats_bucket`, `extended_stats_bucket`, `percentiles_bucket`. ## Where it runs, and why that matters A search fans out to the shards; each shard computes its slice of the aggregation tree; the coordinating node merges those slices in the **reduce phase**. Pipeline aggregations run at the end of that reduce, on the merged result. Three consequences matter in interviews: 1. **They see only what the source aggregation returned.** If a `terms` aggregation returned the top 20 terms, a pipeline aggregation over it knows about exactly those 20 buckets. It cannot pull in a bucket that was never returned, and it inherits any approximation in the source. 2. **They do not reduce work.** Documents are still collected, buckets are still built, and the memory to hold them is still spent. A pipeline aggregation that discards buckets discards them *after* they were paid for. 3. **The cost lands on the coordinating node.** The arithmetic is cheap per bucket, but a very wide bucket set (tens of thousands of buckets) plus several script-based pipeline stages is coordinating-node CPU, not shard CPU. ## What they cannot do Pipeline aggregations take no sub-aggregations — the leaves of the aggregation tree in that sense. They cannot filter documents; only the query does that. They cannot be referenced from a `terms` aggregation's `order` clause, which is precisely why `bucket_sort` exists. And when the value they need is missing for a bucket (an empty bucket where `avg` yields null), the `gap_policy` parameter decides what happens — `skip` is the default, `insert_zeros` substitutes zero. They **can** chain. A `bucket_script` may feed a `derivative`, and a `bucket_sort` may sort on a `bucket_script` output, because each stage's result is just another value hanging off the bucket. ## What they are actually used for The recurring uses are time-series shaping and report-style post-processing: day-over-day change on a `date_histogram`, a running total, a smoothed trend line, a ratio of two metrics per bucket (conversion rate, error rate), a HAVING-style filter that keeps only buckets above a threshold, ordering buckets by a computed value, and a one-line summary such as "the month with the highest sales" or "the average of the monthly totals". If the calculation is a function of already-aggregated numbers, a pipeline aggregation is the right tool; if it is a function of the documents, it is not.

  • Can a pipeline aggregation have sub-aggregations?
    No. Pipeline aggregations accept no sub-aggregations — they consume values and emit a value (or reshape the bucket list), so there is nothing for a child aggregation to run over. They can, however, be chained: a `derivative` can read the output of a `bucket_script`, and a `bucket_sort` can sort on it, because each stage's result is just another value attached to the bucket.
  • Does adding a pipeline aggregation make the search read more documents?
    No. It runs during the reduce phase on results the shards already produced, so no extra document collection happens. The cost is coordinating-node CPU and memory proportional to the number of buckets it walks, plus script compilation and execution for `bucket_script` and `bucket_selector`. The expensive part of such a request is almost always the source bucket aggregation, not the pipeline stage.
  • What happens if a bucket has no value for the metric a pipeline aggregation points at?
    The `gap_policy` parameter decides. The default `skip` treats the bucket as if it did not exist, so no value is emitted there and the calculation continues from the next available value. `insert_zeros` substitutes zero instead, which is right for counts and sums but wrong for averages and rates, where it invents a data point.

A metric aggregation is a formula over raw rows; a pipeline aggregation is a formula over a pivot table's cells — it never looks at the rows again, only at the numbers the pivot already produced.

saying these in an interview costs you the question

  • Says a pipeline aggregation filters or re-reads documents
  • Thinks it runs on each shard during collection
  • Claims it makes the aggregation cheaper by pruning buckets
  • Puts sub-aggregations inside a pipeline aggregation
  • Believes it can see buckets the source aggregation never returned

context