skip to content

How would you stop dashboard aggregations over 90 days of logs from destabilizing a shared Elasticsearch cluster?

level: principalimportance: should knowfreq 32%

answer

  1. Measure before tuning: slow log, breakers, tasks
  2. Old data does not belong on the write path
  3. The same question asked repeatedly should be precomputed
  4. Repeatable requests can be made cacheable
  5. Containment and isolation before capacity

basics

~20 s

Measure which panels cost what, then shrink the input with tiering and retention, pre-aggregate the recurring analyses into summary indices, make the remaining queries cacheable, and bound the blast radius with bucket limits and workload isolation before buying heap.

solid answer

~40 s

Start by measuring rather than tuning: the search slow log, `_nodes/stats/breaker` trip counters and the tasks API tell you which panels are expensive and whether they hurt ingest. Then work a ladder. **Shrink the input** — enforce retention and ILM so 90-day queries touch cold or frozen tiers rather than hot nodes, and require a bounded time range in the UI. **Pre-aggregate** the recurring analyses: a continuous transform writing hourly summaries, or downsampling for time-series data, turns nightly heap crises into a lookup against a tiny index. **Make the rest cacheable** with `size: 0`, rounded date math and longer refresh intervals on read-mostly indices. **Bound the blast radius** with `search.max_buckets`, application-side caps on `terms` size, and dedicated coordinating nodes so dashboard traffic cannot starve ingest. Capacity is the last rung, not the first.

code

json · 15 lines
json
{
  "source": { "index": "logs-*" },
  "dest": { "index": "logs-hourly-summary" },
  "sync": { "time": { "field": "@timestamp", "delay": "60s" } },
  "pivot": {
    "group_by": {
      "service": { "terms": { "field": "service.name" } },
      "hour": { "date_histogram": { "field": "@timestamp", "fixed_interval": "1h" } }
    },
    "aggregations": {
      "events": { "value_count": { "field": "service.name" } },
      "p95_ms": { "percentiles": { "field": "duration_ms", "percents": [95] } }
    }
  }
}

go deeper

for a junior

Recall the basics you control in a query: set size to 0, filter to a narrow time range, and keep facet sizes small so the request stays cheap.

for a middle

Explain the mechanisms behind the advice — why the fetch phase matters, why the request cache needs size: 0, why nested aggregations multiply buckets, and why old indices should refresh less often.

for a senior

Demonstrate the diagnostic and remediation loop on a real cluster: slow log and breaker stats to find the culprits, then tiering, transforms and bucket limits, with node-role separation to protect ingest.

for a principal

Own the tradeoff and the policy: what interactive analytics is worth in hardware, where isolation boundaries go, whether tenants get quotas, and how new dashboards get their cost reviewed before they ship.

## Frame it as a workload problem, not a query problem The failure being described is not one slow query; it is an interactive analytics workload sharing hardware with an ingest workload, where the analytics side has unbounded ability to demand memory. Any answer that starts with a setting has already lost the plot. The order that works is: measure, reduce the input, move work off the query path, make what remains cheap and cacheable, contain the damage the rest can do, and only then talk about capacity. ## Measure first You need to know which panels cost what and how the cost lands on the cluster. - The **search slow log** (`index.search.slowlog.threshold.query.warn` and the fetch equivalents) attributes time at shard level, which is where the expensive work happens. - `GET _nodes/stats/breaker` gives per-breaker `tripped` counters; a rising request-breaker count is dashboard aggregations, a rising fielddata count is a mapping problem. - `GET _tasks?actions=*search*&detailed` shows what is in flight when things go bad, including the request body. - Node-level GC and heap graphs tell you whether the analytics load is what is actually stalling ingest, or whether you are chasing a correlation. Also characterize the *shape*: is it a few enormous ad-hoc queries, or thousands of cheap panels? Those need different fixes, and guessing wrong wastes a quarter. ## Shrink the input Most 90-day aggregation pain is self-inflicted by touching hot-tier data that should have moved on. Time-based indices or data streams with ILM let the old days sit on warm, cold or frozen tiers with fewer replicas and cheaper hardware, so a 90-day query no longer competes with indexing for the same nodes' heap and page cache. Retention that actually deletes data is the cheapest performance work available. On the product side, a bounded default time range and a deliberate friction step for very wide ranges removes far more load than any tuning. "Anyone may drag the range to a year on any panel" is a capacity decision that was made accidentally. ## Move the work off the query path The defining property of a dashboard is that it asks the same questions repeatedly. That is precisely the situation pre-aggregation is for. - A **continuous transform** consumes the raw index and maintains a summary index — per hour, per service, per status — that the dashboard queries instead. The aggregation runs once per bucket, not once per viewer. - **Downsampling** for time-series data streams does the equivalent for metrics, replacing fine-grained points in old intervals with statistical summaries. Elastic's direction of travel in 8.x is toward downsampling for time-series rather than the older rollup APIs, so new work should target it. The result is a dashboard that reads a small index, where even a cache miss is cheap and the aggregations are trivially bounded. ## Make what remains cheap and cacheable Every panel should be `"size": 0`, which skips the fetch phase and is the precondition for the shard request cache. Round date math (`now-24h/h`) so consecutive reloads ask an identical question. Raise `refresh_interval` on indices that are no longer being written, so cache entries survive. Set `track_total_hits: false` where no count is displayed. Individually small wins; collectively they change the load profile of a hundred panels. ## Bound the blast radius Assume something expensive will still get through, and make its failure local. - `search.max_buckets` caps the buckets any one search may produce and fails it deterministically rather than after it has consumed heap. - Cap `terms` `size` in the application or query-building layer; users do not need a hundred-thousand-bucket facet, and no UI renders one. - Keep the circuit breakers at their defaults. They are the containment mechanism; raising them to make errors stop is how a rejected request becomes a node-wide GC stall. - **Isolate the workloads.** Dedicated coordinating nodes absorb the merge-and-reduce phase away from data nodes; dedicated master nodes keep cluster state stable when data nodes struggle. Where the analytics workload is genuinely adversarial to ingest, separate clusters with cross-cluster search is a legitimate answer rather than a defeat. ## Then, capacity More nodes and more heap do help — with the standard caveat that heap beyond roughly 30 GB loses compressed ordinary object pointers, so scaling out beats scaling up, and that half the machine should stay available to the page cache for doc_values and segment reads. Shard sizing matters too: too many small shards makes every aggregation pay per-shard overhead on the coordinating node. ## Governance is part of the answer A principal-level answer names the organizational half: who may create panels, whether tenants have quotas, whether new dashboards get reviewed for cost the way schema changes do, and what the SLO is for interactive analytics versus ingest. Otherwise you fix the cluster and the same problem returns with next quarter's dashboards.

  • Why is aggregating on a runtime field particularly expensive, and when is it acceptable?
    A runtime field has no doc_values, so its script is evaluated for every document the aggregation touches, and it cannot use global ordinals. Cost scales with matched documents and script complexity. It is fine for rare ad-hoc exploration or a schema-on-read patch over a small, filtered slice; anything on a dashboard's hot path should be indexed properly and backfilled.
  • When is a separate cluster the right answer rather than more tuning?
    When the analytics and ingest workloads have genuinely conflicting requirements — different availability targets, unpredictable ad-hoc query cost, or one tenant able to harm another — and no amount of node-role separation gives isolation strong enough. Cross-cluster search preserves a single query surface while putting a hard boundary between the failure domains.
  • What does search.max_buckets protect against that circuit breakers do not?
    It is a deterministic count limit on how many buckets one search may produce, so a bucket-multiplying nested aggregation fails predictably on request shape rather than depending on how much heap happens to be free. Breakers are memory estimates that trip differently under different load; this one fails the same way every time, which makes it usable as a contract with dashboard authors.
  • How do you justify pre-aggregation when the raw data must still be queryable?
    Keep both. The transform's summary index serves the dashboards and the recurring reports; the raw index stays for drill-down, incident investigation and anything ad hoc, ideally on a cheaper tier. Most read volume moves to the cheap path while the expensive path is used rarely and by people who accept its latency.

saying these in an interview costs you the question

  • Reaching for bigger heap or higher breaker limits first
  • Treating it as one slow query rather than a workload conflict
  • Ignoring data tiering and retention entirely
  • Assuming caching alone fixes ad-hoc range changes
  • No plan for who is allowed to create expensive panels

context