skip to content

When an MPP warehouse adds transient clusters as queries queue, which slowness does that fix and which does it not?

level: seniorimportance: should knowfreq 50%

answer

  1. adds concurrency, not horsepower
  2. split waiting-to-start from waiting-to-finish
  3. one query still runs on one cluster
  4. new capacity arrives with empty caches
  5. elasticity multiplies the bill too

basics

~20 s

Adding clusters adds concurrency, so it removes queue wait when many queries arrive at once. It does not make an individual query faster: a single long scan runs on one cluster at the same speed, and new clusters start with cold caches.

solid answer

~50 s

Autoscaling by adding clusters is **horizontal concurrency**, not more power per query. Each query still executes on one cluster, so if the complaint is "this one report takes 30 minutes", extra clusters change nothing — that needs a bigger cluster, better pruning, or a pre-aggregate. What it does fix is the burst case: 200 analysts logging in at 09:00, where elapsed time is dominated by queue wait rather than execution. Decompose the latency first; if wait is small and execution is large, autoscaling is the wrong tool. Even where it fits, know the caveats: a freshly started cluster has cold local caches, so its first queries read from remote storage and run slower than the same query on a warm cluster; there is a scale-up delay before the cluster is usable; and cost scales with the number of clusters, so an unbounded policy plus a runaway query pattern can multiply the bill. Cap the maximum and pair it with per-query guardrails.

code

text · 8 lines
text
burst at 09:00, 200 dashboard queries
  before: queue_wait 45s + execution 2s   = 47s
  after adding clusters: queue_wait 1s + execution 2s (warm) = 3s
                          queue_wait 1s + execution 6s (cold cache) = 7s

single nightly report, runs alone
  before: queue_wait 0s + execution 31m = 31m
  after adding clusters: queue_wait 0s + execution 31m = 31m   <- no change

go deeper

for a junior

Know the distinction between a bigger cluster, which speeds up one query, and more clusters, which let more queries run at the same time.

for a middle

Be able to decompose elapsed time into queue wait and execution and say which of the two extra clusters actually addresses, plus why a fresh cluster starts slower than a warm one.

for a senior

Show judgment on fit: read-mostly bursty workloads of many small queries benefit; a few dominant large queries or a write-heavy pipeline do not. Name the guardrails you pair with elasticity.

for a principal

Own the cost model — the maximum cluster count you are willing to fund, what triggers it, and the policy that stops elasticity from amplifying a defective workload into a billing event.

## Two different kinds of "more compute" Warehouses scale in two orthogonal directions and conflating them is the classic error. - **Scale up (bigger cluster)**: more nodes, cores and memory available to *one* query. This shortens a single query's execution time, up to the point where the query stops being parallelisable or hits skew. - **Scale out (more clusters)**: additional identical compute units, each running its own queries against the same shared storage. This raises how many queries can run **simultaneously**. Each individual query still executes within one cluster and is no faster than it was. Autoscaling that spins up transient clusters when queries start queueing is the second kind. So the first question in any diagnosis is: **is the user waiting to start, or waiting to finish?** ## Decompose before you scale Elapsed time from the user's perspective is roughly: ``` elapsed = queue_wait + compile + execution (+ result fetch) ``` - **Queue wait dominant** — many queries, limited slots. Adding clusters directly attacks this and works well. - **Execution dominant** — the query itself is expensive. Adding clusters does nothing. Fix the query: prune more, cluster the data on the filter column, join in a better order, or serve from a materialised aggregate. Or scale the cluster *up* so the query gets more parallelism. - **Compile dominant** — usually enormous generated SQL or metadata pressure; scaling is irrelevant. A useful signal: if queue wait rises and falls with the number of active users but execution time is flat, the workload is concurrency-bound and autoscaling is the right instrument. ## What autoscaling costs you **Cold caches.** Warehouses lean heavily on local SSD or memory caching of remote storage. A newly started cluster has none of it. The same query that takes 3 seconds on the warm original cluster may take noticeably longer on the transient one because it must fetch from remote object storage. This produces confusing support tickets: two users running the same dashboard at the same second get visibly different latencies, purely because one landed on a warm cluster and one on a cold one. Result caches that are cluster-independent mitigate this for repeated identical queries; local data caches generally do not transfer. **Start-up delay.** Provisioning is not instantaneous. A burst shorter than the provisioning time is over before the extra capacity arrives, so autoscaling helps sustained bursts more than spiky ones. **Cost.** Concurrency scaling multiplies compute spend by the number of active clusters. Unbounded scaling plus a workload that generates unbounded queries — a dashboard set to auto-refresh every 10 seconds across 500 tabs, or a retry loop — turns a performance feature into a billing incident. Always set a maximum cluster count and pair autoscaling with per-query guardrails, so scaling absorbs legitimate demand rather than amplifying a defect. **Not always applicable to every statement.** Platforms typically route only certain kinds of work — commonly read-only queries — to transient clusters, since concurrent writers introduce coordination the feature is not designed to arbitrate. Do not assume a heavy write or maintenance operation gets the same elasticity as a read. ## Where it fits best Autoscaling extra clusters is close to ideal when the workload is: - **read-mostly**, so routing to another cluster raises no write-coordination question; - made of **many small-to-medium queries** rather than a few huge ones, so per-query execution time was never the constraint; - **bursty on a human schedule** — the 09:00 login wave, the Monday-morning reporting surge — where paying for peak capacity all day would be wasteful. It is a poor fit when a handful of very large queries dominate, when the workload is write-heavy, or when the underlying problem is one badly written query that nobody has fixed. ## Combine, don't substitute The mature configuration uses all three levers together: admission control to keep each cluster near the top of its throughput curve, separate pools so classes of work cannot interfere, and autoscaling within a pool to absorb legitimate bursts — with an upper bound and per-query limits so the elasticity never becomes the failure mode.

  • Why can two users running the identical dashboard at the same moment see very different latencies under autoscaling?
    They are likely served by different clusters. The original cluster has warm local caches holding recently read data, while a just-started transient cluster must fetch the same data from remote storage. Same SQL, same plan, different cache state. A cluster-independent result cache hides this for byte-identical repeated queries, but any query that actually scans will show the gap until the new cluster warms up.
  • What guardrail would you pair with autoscaling to keep it from becoming a cost incident?
    Two. First, a hard cap on the number of concurrent clusters, so demand cannot translate into unbounded spend. Second, per-query limits — a runtime timeout and a scanned-bytes ceiling — so a defective query or a runaway refresh loop is stopped rather than being given more capacity. Add a spend monitor that alerts, and optionally suspends the pool, at a threshold, so the failure mode is a paused workload rather than a surprise invoice.
  • When is scaling the cluster up the right answer instead of adding clusters?
    When execution time dominates and the query is genuinely parallelisable: a large scan, sort or join that can use more cores and more memory. Scaling up shortens that one query and reduces spilling by enlarging the memory available for its grant. It stops helping when the query is limited by skew, by a serial stage, or by a single hot key — in that case more hardware idles while one worker does the work.

Opening more checkout lanes clears a queue of shoppers quickly, but it does nothing for the one shopper whose trolley takes twenty minutes to scan. That shopper needs a faster scanner, not another lane.

saying these in an interview costs you the question

  • Expects extra clusters to shorten a single long-running query
  • Ignores cold caches on newly started clusters
  • Assumes autoscaling is free or cost-neutral
  • Never bounds the maximum number of clusters
  • Treats autoscaling as a substitute for fixing a bad query

context