For a system that runs both short transactional queries and large reporting queries against the same relational engine, how would you decide how much working memory to allow hash join operators, and what are the consequences of setting that budget too high or too low?
answer
- grant is per operator, per worker, per query
- demand = grant × ops × workers × concurrency
- operator memory is stolen from the buffer cache
- conservative global default + scoped exception
- spilling by design for batch is fine
basics
~20 sBudget from total memory minus the buffer cache, divided by realistic peak concurrency and the number of memory-hungry operators per query — not per query in isolation. Too low means constant spilling and IO; too high means a few concurrent reports exhaust memory, evict the cache, or crash the process.
solid answer
~60 sTreat per-operator memory as a **global budget divided by concurrency**, not a per-query setting. The rough model: usable memory = RAM − OS − buffer cache − connection overhead. Then note that the grant is per *operator*, not per query: one report can have several concurrent hash and sort operators, and each parallel worker may take its own grant. So worst-case demand is roughly `grant × operators_per_query × parallel_workers × concurrent_queries` — a number that surprises people by an order of magnitude. Set it **too low** and large joins spill constantly: extra IO passes, temp-storage bandwidth contention, sharp latency cliffs on reports. Set it **too high** and a handful of simultaneous reports can consume all memory, starve the buffer cache (turning cached reads back into IO for everyone, including OLTP), or get processes killed. The practical answer is workload separation: a conservative global default sized for OLTP concurrency, a raised setting scoped to the reporting role, session, or job, and — better still — reports run on a replica or a separate instance so their memory behaviour cannot affect transactional latency. Then verify with spill metrics rather than theory.
go deeper
Know that hash joins need working memory, that too little causes spilling to disk, and that the setting applies to many operators at once rather than to the machine as a whole.
Do the arithmetic: grant times operators times workers times concurrent queries, and explain that memory given to operators is taken from the cache.
Recommend a conservative default with scoped exceptions for reporting sessions, back it with spill and cache metrics, and fix query shape and estimates before touching the number.
Frame it as workload isolation and admission control — separate reporting onto a replica, cap heavy-query concurrency, provision temp storage for deliberate spilling, and define which workloads are allowed to degrade under contention.
## Why this is a judgment question There is no correct number, because the setting trades three things against each other: **per-query speed** (avoid spilling), **concurrency** (how many queries can run without exhausting memory), and **cache effectiveness** (memory not given to operators stays in the buffer cache, benefiting everything). An interviewer asking this wants to see whether you reason from a budget and a workload model rather than from a remembered default. ## The accounting that people get wrong Start with the machine: total RAM, minus OS and filesystem needs, minus the buffer cache (usually the largest single allocation and the highest-leverage one), minus per-connection overhead, leaves the pool available for query working memory. Now the multiplier that catches teams out — the grant is generally **per operator**, not per query: - one query can contain several memory-hungry operators at once (two hash joins and a sort in one plan is ordinary); - if the engine parallelizes, **each worker** may take its own grant for its own operator; - and all of that multiplies by the number of **concurrent** such queries. So worst-case demand ≈ `grant × operators × workers × concurrent_heavy_queries`. A setting that looks modest for one query can, at a realistic peak, demand many times the machine's memory. Any sizing that ignores this multiplier is wrong by construction. ## Consequences of setting it too low - **Chronic spilling.** Hash joins fall back to partitioning: extra write and read passes over both inputs. For a big report this can be a several-fold slowdown. - **Temporary storage becomes a shared bottleneck.** Many spilling queries contend for the same temp devices and bandwidth, so the slowdown is superlinear in concurrency. - **Cliff-edge behaviour.** The transition from fitting to spilling is a step change, so a modest data-growth increment turns a fast query slow overnight — one of the most common "nothing changed and it broke" incidents. - **Extra recursion.** A budget far below the build size forces more partitions and possibly recursive repartitioning, each level adding another IO pass. ## Consequences of setting it too high - **Memory exhaustion under concurrency.** The multiplier above bites: three reports at once, each with two hash joins and four workers, can demand 24 grants. - **Buffer cache starvation.** Memory handed to operators is memory not caching data pages. Ironically, a huge working-memory setting can make everything slower by evicting the cache the OLTP workload depends on. - **Blast radius crossing workloads.** Reports and transactions share one memory pool, so an analyst's ad-hoc query becomes a customer-facing latency incident. - **Hard failures.** Depending on the platform, over-commitment ends in allocation errors or the OS killing the database process — an availability event, not a performance one. ## A defensible method 1. **Characterise the workload.** How many concurrent heavy queries genuinely occur at peak? What are the real build-side sizes of the top reports? This is measurable, not guessable. 2. **Set a conservative global default** sized for the OLTP concurrency level. The majority of statements need almost no working memory; the default should serve them, not the outliers. 3. **Scope the exception.** Raise the setting per role, per session, or per batch job for the queries that genuinely need it, so the exposure is bounded to a known small number of sessions. 4. **Separate the workloads physically where possible.** Running reports on a replica or a dedicated instance is the strongest answer: it removes the memory interaction entirely and, as a bonus, removes the cache-pollution and long-transaction effects that reporting queries have on a transactional primary. 5. **Bound concurrency, not only size.** A queue or admission limit on heavy queries controls the multiplier directly and is often more effective than tuning the per-operator number. 6. **Measure and iterate.** Track spill volume and frequency, temp-storage usage, cache hit ratio, and peak memory. Tune toward "the important reports do not spill at realistic concurrency" rather than "nothing ever spills". ## Accept spilling deliberately A critical piece of judgment: **spilling is not a failure mode to eliminate at any cost.** Grace/hybrid partitioning exists precisely so joins can exceed memory, and for genuinely huge batch joins a stable spilled plan is the correct, predictable design. Provisioning enough memory to hold a 200 GB build side in RAM is not a sensible goal. The goal is that *latency-sensitive* work does not spill unexpectedly, and that batch work spills predictably on storage provisioned for it. ## Cheaper levers before touching the number Often the right fix is not the budget at all: - **Narrow the build side** — projecting only needed columns shrinks the hash table proportionally, and wide `SELECT *` reports are a common cause of overflow. - **Filter and pre-aggregate earlier**, so fewer rows are hashed. - **Fix estimates** — statistics and multicolumn statistics, so the engine sizes partitions sensibly rather than repartitioning at runtime. - **Reconsider the plan shape** — sometimes a merge join over already-ordered input is cheaper and has bounded memory behaviour. ## What a strong answer sounds like "Budget from the machine down, divide by realistic peak concurrency times operators per query times parallel workers, keep the global default conservative for OLTP, raise it only for scoped reporting sessions, prefer moving reports off the primary entirely, cap heavy-query concurrency, and validate with spill and cache metrics — accepting that large batch joins should spill by design."
- Why can raising the per-operator memory setting make the whole system slower?Memory allocated to query operators is memory unavailable to the buffer cache, so a large setting can shrink the cache and turn previously cached reads into physical IO for every workload including OLTP. Under concurrency the effect compounds, because the grant applies per operator and per parallel worker rather than per query. In the worst case allocation failures or an OS-level kill turn a tuning choice into an availability incident.
- How would you give big reports the memory they need without endangering transactional latency?Scope the exception rather than raising the global default: set a higher working-memory value per reporting role, session, or batch job so only a known, small number of sessions can claim large grants. Better still, run reporting on a replica or separate instance so its memory, cache, and long-transaction behaviour cannot touch the primary. Pair either approach with an admission limit on concurrent heavy queries, since concurrency is the multiplier that actually causes exhaustion.
- Should the goal be that no query ever spills?No. Partitioned spilling is the designed mechanism for joining inputs larger than memory, and provisioning RAM to hold an enormous build side is neither possible nor economical. The goal is that latency-sensitive queries do not spill unexpectedly, and that batch work spills predictably onto temporary storage provisioned for that throughput. Unexpected spills are usually an estimation or query-shape problem rather than a budget problem.
- Before changing the memory setting, what query-level fixes would you look for?Narrow the build side by projecting only the columns actually needed, since a wide select inflates the hash table in direct proportion to row width. Push filters down and pre-aggregate so fewer rows are hashed at all. Fix cardinality estimates with refreshed and multicolumn statistics so the engine sizes partitions correctly, and check whether a merge join over already-ordered inputs would give more predictable memory behaviour.
Handing out workbench space in a shared workshop: give each job a huge bench and two jobs fill the room; give everyone a tiny bench and every job spends its time shuttling parts to the storeroom. You size the bench from the room and the expected number of simultaneous jobs — and you put the messy long jobs in a different room.
saying these in an interview costs you the question
- Sizing the setting for one query in isolation, ignoring concurrency and parallel workers
- Treating the value as per query rather than per operator
- Forgetting that operator memory competes with the buffer cache
- Aiming to eliminate all spilling regardless of data size
- Raising the global default to satisfy a handful of reports