skip to content

MPP & Warehouse Architecture

How analytical systems spread a single query across many machines, and where the data physically sits relative to the compute working on it. Interviewers use this axis to see whether I can reason about scaling, shuffles and cost rather than only writing SQL.

on this pageshow

questions

24

In a cloud data warehouse, what does separating storage from compute actually mean?

level: juniorimportance: must knowfreq 76%

answer

  1. one copy of data, many workers
  2. nodes no longer own their slice
  3. data on object storage, nodes are stateless
  4. turn compute off, data stays
  5. resize CPU without moving bytes

basics

~20 s

Table data lives as files in a shared object store that every node can read over the network, and compute clusters are stateless workers that read it. Either side can be resized, paused or duplicated without touching the other.

solid answer

~50 s

In a classic shared-nothing warehouse each node owns local disks holding a slice of every table, so storage capacity, CPU and node count are welded together. In a separated architecture the table data is written once as immutable columnar files into a shared object store, and a metadata service records which files make up the current version of each table. A compute cluster is a group of stateless VMs that ask metadata what to read, pull byte ranges over the network, and cache what they read on local SSD and in memory. Nothing authoritative lives on the cluster, so you can suspend it, resize it, or start a second one over the same data without copying or redistributing anything. Storage grows on its own cost curve; compute is bought only for the hours you actually query.

go deeper

for a junior

Be ready to state plainly that table data sits in object storage while query nodes are temporary workers, and that shutting compute down does not lose data.

for a middle

Explain the mechanics: immutable files plus a metadata service defining the current version, stateless workers with local SSD caches, and why resizing no longer requires redistributing data.

for a senior

Show the trade-off you actually operate: a network read path, cold caches after resume, and a metadata tier that every query and commit depends on. Say what you tune to buy locality back.

for a principal

Own the economics. Argue when two independent cost curves beat one coupled one, what elasticity is genuinely worth for your duty cycle, and where the coupling reappears as a metadata or cache bottleneck.

## The architecture it replaced A classic massively parallel warehouse is shared-nothing: every compute node has its own attached disks, and every table is split across those disks so each node owns a slice. The node that owns the bytes is the node that scans them, which makes reads local and fast. The price is that three separate decisions get welded into one number — the node count. How much data you can hold, how much CPU you can throw at a query, and how much you pay per hour are all the same knob. If your data doubles you buy nodes even though query volume is flat. If Monday morning needs four times the CPU you buy nodes and then wait while the cluster redistributes data onto them. And you cannot switch the cluster off overnight, because switching it off takes the data offline. ## What the separated design looks like In a separated architecture there are three tiers rather than one. **Storage.** Table data is written as immutable columnar files into a shared object store — S3, Google Cloud Storage, Azure Blob Storage or an equivalent. The object store handles durability and replication. Every compute node can read every file over the network; no node owns anything. **Metadata.** A separate transactional service records which set of files constitutes the current version of each table, plus per-file statistics used to skip files. A write does not modify existing files; it adds new ones and commits a new file list. Because the commit is a single metadata operation, readers see either the old set or the new set and never a half-written table. **Compute.** A cluster is a group of stateless worker VMs. To run a query it asks metadata which files matter, issues parallel ranged reads against the object store, and decodes the columns it needs. Whatever it reads it caches — in memory and on local NVMe SSD — purely as an optimization. Killing the cluster loses only cache. ## What the split buys you - **Independent scaling.** Storage grows without buying CPU, and CPU scales for a busy hour without moving a byte. Resizing is provisioning VMs, not rebalancing data. - **Suspend and resume.** Idle compute can be stopped. Storage keeps costing; compute stops costing. This is what makes per-second or per-query billing possible at all. - **Many clusters, one copy.** Several independently sized clusters can read the same tables at the same time, because none of them owns the data. Previously a second workload meant a second copy and a pipeline to keep it fresh. - **Durability is delegated.** The object store's replication is the durability story, so node loss is a compute event rather than a data event. - **Elasticity of shape, not just size.** The same dataset can be queried by a small cluster for a dashboard and a large one for a nightly rebuild. ## What it costs The read path now crosses a network. Object storage has far higher per-request latency than a local NVMe drive and rewards a few large sequential reads over many small ones, so engines lean hard on large columnar files, column pruning, file-level statistics, deep read-ahead and aggressive local caching. Right after a cluster starts, all of that caching is empty, so the first queries run at remote-read speed — the familiar cold-start penalty. The metadata service also becomes a shared stateful dependency that every query start and every commit touches, so it is a real scaling and availability consideration rather than an implementation detail. ## What it does not change Separation is about where bytes live, not about how queries execute. Joins still redistribute rows between nodes across the network. Aggregations still shuffle. A query that reads a whole fact table still reads a whole fact table — and now pays network for it, which is why layout and pruning matter at least as much as before, not less. Nor does it make compute free: a stateless cluster that sits running with no queries burns exactly as much money as a busy one. ## How to talk about it The honest summary is that separation converts a capacity problem into two independent cost curves and buys elasticity, and it pays for that with a network hop and cold caches. Engines then spend enormous effort — caching tiers, statistics, file sizing, prefetch — buying back the locality they gave up.

  • If compute nodes hold no authoritative data, why do these systems still attach fast local SSDs to them?
    Purely for caching. Object storage has high per-request latency, so the engine keeps recently read column chunks on local NVMe and in memory. A warm cluster serves most of a repeated scan locally and only fetches misses remotely, which is why the same query gets markedly faster on its second run.
  • Does separating storage and compute mean data layout and partitioning stop mattering?
    No — it usually matters more. A scan that fails to prune now pays network transfer and, under per-byte-scanned billing, pays directly in dollars. Good file layout, sensible clustering and selective filters are what keep a remote read path fast; the separation removes the capacity coupling, not the physics of reading bytes.
  • What is the biggest new failure domain this architecture introduces?
    The metadata and catalog service. Every query start resolves the current file list through it and every commit is a write to it, so its availability and latency bound the whole system. Object storage outages matter too, but the metadata tier is the piece the vendor had to build and scale itself.

A shared-nothing cluster is a library where each librarian keeps their own shelf behind their desk; sending one home takes those books with them. Separation puts every book in one central stack and hires readers by the hour.

saying these in an interview costs you the question

  • Claims queries get faster because storage is now remote
  • Says the cluster still holds a permanent copy of each table
  • Thinks separation removes shuffle or network cost from joins
  • Believes suspended compute still bills, or that storage stops billing
  • Assumes you must add nodes as data volume grows

context

open as a page

What is the difference between a data warehouse and a data lake for analytical data?

level: juniorimportance: must knowfreq 76%

basics

~20 s

A data warehouse stores data in an engine-managed, usually proprietary format with an enforced schema and integrated compute. A data lake stores raw files in object storage that many engines can read, with structure applied only at query time.

open as a page

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

level: middleimportance: must knowfreq 68%

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.

open as a page

In an MPP query plan, what is an exchange operator and when must the planner insert one?

level: middleimportance: must knowfreq 74%

basics

~20 s

An exchange moves rows between compute nodes over the network. The planner inserts one wherever an operator needs matching key values on the same node and the current data placement does not already guarantee that — joins, grouping, global sorts, and the final gather.

open as a page

In a cloud warehouse, why is the first query after a suspended compute cluster resumes slower than later runs?

level: middleimportance: must knowfreq 68%

basics

~20 s

Resuming provisions fresh worker VMs whose memory and local SSD caches are empty, so the first query reads everything remotely from object storage. Subsequent runs hit warm local caches and skip most of that network traffic.

open as a page

What does the lakehouse architecture add on top of raw files sitting in a data lake?

level: middleimportance: must knowfreq 64%

basics

~20 s

A lakehouse adds an open table-format metadata layer over lake files, so a directory of files behaves like a table: atomic commits, consistent snapshots for readers, schema enforcement and evolution, row-level updates and deletes, and statistics for pruning — while staying readable by many engines.

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

One MPP node runs at 100% while the rest idle during a large GROUP BY — what is happening and how do you fix it?

level: seniorimportance: must knowfreq 66%

basics

~20 s

The shuffle key is skewed: a few values, often NULL or a placeholder, own most of the rows, so one node receives that whole bucket. Runtime equals the slowest node. Fix by isolating or salting the hot values, or by pre-aggregating before the shuffle.

open as a page

In a shared-nothing MPP warehouse, how is one table's data divided across the cluster?

level: juniorimportance: should knowfreq 55%

basics

~20 s

Each compute node owns an exclusive subset of the table's rows and can read only that subset. Rows are placed by hashing a column, by round-robin, or by replicating the whole table everywhere. Every node then runs the same plan over its own share.

open as a page

How does an MPP planner choose between broadcasting one join input and hash-redistributing both?

level: middleimportance: should knowfreq 63%

basics

~20 s

It compares network cost. Broadcasting sends the smaller input's bytes to every node, so cost scales with node count; redistributing sends each row of both inputs once. Broadcast wins when one side is small enough to fit in every node's memory.

open as a page

In a cloud warehouse, how does a result cache differ from a compute cluster's local data cache?

level: middleimportance: should knowfreq 55%

basics

~20 s

A result cache stores finished query results in the shared service layer and can answer a repeat query without running compute at all. A local data cache stores raw column chunks on a cluster's own SSDs and still requires the cluster to execute the query.

open as a page

What do schema-on-write and schema-on-read mean for an analytical data platform?

level: middleimportance: should knowfreq 58%

basics

~20 s

Schema-on-write validates structure when data is loaded, so anything stored already conforms and bad rows are rejected at the source. Schema-on-read stores data as it arrives and interprets it at query time, so ingestion is cheap but each consumer bears the parsing and quality risk.

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

Why does an MPP query spill to disk after a shuffle even when total cluster memory looks sufficient?

level: seniorimportance: should knowfreq 48%

basics

~20 s

Cluster memory is not one pool. Each worker gets a fraction of one node's RAM, divided again among concurrent queries, and a shuffle can hand one worker far more than its even share. That worker spills alone while the cluster looks half empty.

open as a page

When does per-byte-scanned billing beat paying by the second for a provisioned warehouse cluster?

level: seniorimportance: should knowfreq 60%

basics

~20 s

Per-byte-scanned billing wins on spiky, low-duty-cycle workloads whose queries prune well, because you pay nothing between queries. Provisioned compute wins when a cluster stays busy, when queries repeatedly scan the same large data, and when caching amortizes across many statements.

open as a page

When should analytical data stay in external tables on object storage instead of being loaded into the engine?

level: seniorimportance: should knowfreq 54%

basics

~20 s

Keep data external when it is queried rarely, is large and cheap to retain, is still raw or evolving, or must stay readable by other engines. Load it when queries are frequent, latency-sensitive or highly concurrent, because loading buys layout control, statistics, clustering and caching.

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

Doubling an MPP cluster's nodes barely speeds up a shuffle-heavy query — how do you diagnose it and what do you change?

level: principalimportance: should knowfreq 40%

basics

~20 s

Scale-out only shrinks the local work. Redistribution volume is roughly constant in node count, broadcast volume grows with it, skewed keys stay on one node, and the final gather stays serial. Attribute the runtime per stage before buying capacity.

open as a page

In a warehouse with storage and compute separated, stored data triples yearly while query volume stays flat — how do you plan cost?

level: principalimportance: should knowfreq 44%

basics

~20 s

Forecast the two curves separately: storage cost tracks retained bytes, compute cost tracks queries and duty cycle. The trap is second-order coupling — unpruned queries and a working set that outgrows local caches turn storage growth into compute growth.

open as a page

How would you choose between a managed warehouse's own storage and open table formats in your object storage?

level: principalimportance: should knowfreq 40%

basics

~20 s

Decide on how many engines must read the data, how much governance and performance you need out of the box, and how much platform work you can staff. Warehouse-managed storage buys integration and predictability; open formats buy engine choice and exit options, paid for in operational ownership.

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

What are the trade-offs of federated queries from an analytical engine into an operational database?

level: middleimportance: nice to knowfreq 32%

basics

~20 s

Federation lets an analytical engine query a live operational database directly, avoiding a pipeline and giving current data. It costs you: the analytical engine cannot prune or parallelise inside the remote system, large pulls cross the network, and heavy queries add load to a production database.

open as a page

If warehouse compute is stateless and object storage is passive, where do table metadata and transactions live?

level: seniorimportance: nice to knowfreq 36%

basics

~20 s

In a separate transactional metadata service. It records which immutable files make up the current version of each table, their statistics, and the commit order, so a write is a metadata swap and every compute cluster sees one consistent view.

open as a page