skip to content

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