In an MPP analytical warehouse, why can capping the number of concurrently running queries increase total throughput?
answer
- resources per query are finite
- more in flight is not more done
- watch what happens to memory grants
- queries start spilling to disk
- waiting beats thrashing
basics
~20 sA 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 sAn 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
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.
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.
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.
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