skip to content

Concurrency & Workload Isolation

One warehouse serving dashboards, ad-hoc analysts and ETL at once needs admission control, priorities, or a separate compute pool per workload. Interviewers raise it when probing how I would stop a runaway query from starving the BI dashboard.

on this pageshow

questions

6

In an MPP analytical warehouse, why can capping the number of concurrently running queries increase total throughput?

level: middleimportance: must knowfreq 68%

answer

  1. resources per query are finite
  2. more in flight is not more done
  3. watch what happens to memory grants
  4. queries start spilling to disk
  5. waiting beats thrashing

basics

~20 s

A cluster has a fixed budget of memory, CPU and I/O. Past a point, extra concurrency splits that budget so thinly that queries spill to disk and contend for the same resources, so queueing the excess finishes more work per hour.

solid answer

~50 s

An analytical cluster has a fixed budget of memory, CPU, disk and network bandwidth. Admission control decides how large a share one query may hold and how many queries may hold a share at once. Below the saturation point, more concurrency genuinely raises throughput. Above it, each query's memory grant shrinks until hash tables and sorts no longer fit in RAM, so queries spill to disk, re-read the spilled partitions, and then fight each other for that same I/O and network — everything slows at once and the cluster completes *less* total work. Holding the excess in a queue keeps the running set fast: a query that waits 20 seconds and then runs in 10 finishes sooner than fifty queries all crawling. The setting is a throughput-versus-wait trade: too high and the cluster thrashes, too low and it idles while work waits in line.

go deeper

for a junior

Know that a warehouse runs only so many queries at once and holds the rest in a queue, and that this is deliberate rather than a bug.

for a middle

Be ready to explain the mechanics: fixed cluster memory, a per-query grant, spilling to disk when the grant is too small, and the throughput curve that rises then falls as concurrency grows.

for a senior

Expect to justify a specific limit from measurements — representative query footprint, queue wait versus resource utilisation — and to argue for splitting one limit into per-class limits for a mixed workload.

for a principal

Own the trade the number encodes: predictable latency for interactive users against utilisation of paid compute, and the point at which the honest answer is more capacity or a separate pool rather than more tuning.

## The fixed budget An MPP warehouse is a set of nodes with a fixed amount of RAM, CPU cores, local disk bandwidth and interconnect bandwidth. A single analytical query is not a small consumer: a hash join builds a hash table sized by one input, an aggregation builds a hash table sized by the number of groups, a sort needs a buffer sized by the data being sorted, and a shuffle moves rows across the network. All of that has to fit somewhere. Admission control is the component that decides **which queries are allowed to start running now**, and implicitly **how much of the budget each one may take**. Most engines express this in one of three ways, or a mix: - **Slot-based**: the cluster is divided into N concurrency slots; a query occupies one (or several) while it runs. - **Memory-based**: each query is granted a memory budget at admission time, and a query is admitted only if its grant is available. - **Cost-based**: the optimizer's estimate (bytes to scan, expected rows) routes the query into a class with its own limits. ## What happens when nothing is capped Suppose the cluster can comfortably keep eight large queries resident. Admit fifty and several things go wrong at once: 1. **Memory grants shrink.** Each query gets a fraction of what it needs. Hash tables and sorts no longer fit, so operators **spill** — write partitions to local disk and read them back. A spilling hash join can easily do several times the I/O of an in-memory one. 2. **The spill traffic collides.** Fifty queries spilling simultaneously saturate exactly the local disk and network bandwidth the scans also need. Spill makes queries slower, and slower queries stay resident longer, which admits still more pressure — a feedback loop. 3. **CPU is oversubscribed.** Vectorized engines run wide; when runnable threads far exceed cores, context switching and cache-line eviction eat real throughput, because each query's working set keeps getting evicted from CPU cache by the others. 4. **Everything ages together.** The result is not "most queries fast, some slow" — it is *every* query slow, including the trivial ones, and users experience the cluster as broken. Throughput as a function of concurrency is a curve that rises, flattens, then **falls**. Admission control's job is to keep the system near the top of that curve rather than to the right of it. ## Why a queue is the cheap answer Queueing looks like waste — a query sitting idle while the cluster works — but the arithmetic favours it. If eight concurrent queries each finish in 10 seconds, the cluster retires 48 queries per minute; the 49th waits a few seconds and still finishes quickly. If fifty run at once and each takes two minutes because of spill, the cluster retires 25 per minute and *nobody* got a fast answer. Total latency for the whole batch is lower with the queue, and the tail is dramatically better. Queueing also makes behaviour predictable. A bounded running set means the memory grant per query is stable, so a query that took 10 seconds yesterday takes roughly 10 seconds today. Unbounded concurrency makes runtime a function of who else happened to log in. ## Where the limit should sit There is no universal number; it depends on the shape of the workload. Practical method: - Measure the **memory footprint of a representative query** and divide the usable cluster memory by it — that is the order of magnitude for the concurrent count of *that* class. - Separate classes. Tiny dashboard queries touching a few pruned files need almost nothing and can run at high concurrency; a nightly join of two fact tables needs a large grant and must run at low concurrency. One global number for both is always wrong for one of them. - Watch two signals together. **Queue wait high while CPU and memory sit idle** means the limit is too low. **Little queue wait but heavy spilling and degrading runtimes** means it is too high. ## What admission control does not fix It does not make an individual query faster — a badly written 30-minute scan still takes 30 minutes, it just no longer drags everyone down with it. It does not fix skew, missing pruning, or a bad join order. And it does not create capacity: if the arrival rate genuinely exceeds what the cluster can retire, queueing converts a thrashing cluster into an ever-growing backlog. At that point the answer is more compute, a separate pool for the offending workload, or fewer/cheaper queries — the queue only tells you which.

  • How would you tell whether the concurrency limit is set too low rather than too high?
    Look at queue wait against resource utilisation. If queries pile up in the queue while CPU, memory and I/O on the cluster sit well below saturation, the limit is too low and you are idling capacity. If queue wait is small but running queries spill heavily and their runtimes drift upward as concurrency rises, the limit is too high. The healthy state is modest queue wait with resources near, but not at, saturation.
  • Why is one global concurrency number usually wrong for a mixed workload?
    Because memory demand per query varies by orders of magnitude. A pruned dashboard lookup may need megabytes; a fact-to-fact join may need many gigabytes for its hash tables. A limit tuned for the big queries starves the cluster of the cheap ones it could run dozens of; a limit tuned for the cheap ones lets several large queries in at once and they all spill. The fix is multiple admission classes with their own limits and grants.
  • What does a query's memory grant have to do with spilling?
    The grant is the ceiling an operator may use before it must write intermediate data to disk. If a hash join's build side exceeds the grant, the engine partitions it and spills, then re-reads each partition — turning a memory-speed operation into a disk-speed one. Raising concurrency shrinks grants, which is the usual reason a query that ran fine alone spills when the cluster is busy.

A motorway carries the most cars per hour at a moderate density. Let everyone on at once and it jams — total cars delivered per hour drops. A ramp meter that holds cars briefly moves more traffic than an open ramp.

saying these in an interview costs you the question

  • Claims more concurrent queries always means more total throughput
  • Treats any queue wait as a misconfiguration to be removed
  • Confuses time spent queueing with time spent executing
  • Assumes adding nodes makes admission control unnecessary
  • Thinks memory is handed out per query with no ceiling

context

open as a page

Dashboard queries on a shared analytical warehouse jump from 2s to 40s whenever the nightly ETL job runs. How would you isolate the two workloads?

level: seniorimportance: must knowfreq 72%

basics

~20 s

First confirm the slowdown is queue wait or resource contention, not a plan change. Then move ETL onto its own compute pool over the same shared storage — the only fix that removes interference outright. Priorities and memory caps reduce it but never eliminate it.

open as a page

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%

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.

open as a page

What guardrails would you put on an analytical warehouse to stop one runaway query from starving everything else?

level: seniorimportance: should knowfreq 55%

basics

~20 s

Layer them: an admission-time check on estimated bytes scanned, a per-class statement timeout, caps on memory grant and spill, a result-size limit, and a spend monitor that suspends compute. Prefer demoting a query to killing it where possible.

open as a page

For a company-wide analytical warehouse, how do you choose between one compute pool per team and shared pools with priority classes?

level: principalimportance: should knowfreq 40%

basics

~20 s

Split by workload class and SLA, not by org chart. Give guaranteed isolation only where a latency or freshness commitment justifies its idle cost; let everything else share pools with priorities, chargeback by usage, and revisit as workloads change.

open as a page

In a single FIFO admission queue on an analytical warehouse, why does one long query delay hundreds of short ones?

level: middleimportance: nice to knowfreq 38%

basics

~20 s

Because a FIFO queue serves in arrival order: a 20-minute query holding a slot blocks everything behind it, however cheap. This is head-of-line blocking, and the fix is separate admission classes with reserved capacity for short queries.

open as a page