skip to content

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

level: juniorimportance: should knowfreq 55%

answer

  1. no node can read another node's data
  2. a rule decides which node owns a row
  3. hash, round-robin, or copy-everywhere
  4. cores subdivide a node's share further
  5. the same plan runs on every share

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.

solid answer

~50 s

A shared-nothing cluster gives every compute node its own CPU, memory and storage, with no shared buffer pool and no ability to read another node's data directly. A table is split across the nodes by a placement rule: **hash** on a chosen column, so equal key values land together; **round-robin/random**, which spreads evenly but co-locates nothing; or **replicated**, where each node holds a full copy of a small table. Each node further subdivides its share into parallel units — one per core or per storage volume — so a 10-node cluster can have 80 independent scan streams. The optimizer compiles one plan and every unit runs the identical fragment over its own rows. A coordinator node distributes those fragments and gathers the final result. Any operation that needs rows from another node must move them over the network.

go deeper

for a junior

Be ready to say that each node owns a disjoint slice of the rows, cannot read another node's slice directly, and runs the same query fragment over its own slice.

for a middle

Explain the three placement rules — hash on a key, round-robin, replicated — and what each one implies for a later join. Mention that a node's share is subdivided again across cores.

for a senior

Show that you reason about placement before performance: which key the table is distributed on, how evenly it spreads, and which queries that choice makes local versus network-bound.

for a principal

Own the fleet-level consequence: placement decisions are long-lived, expensive to change on large tables, and constrain every future query shape, so they belong in schema review rather than in query tuning.

## What "shared-nothing" means In a shared-nothing massively parallel processing (MPP) cluster, every compute node is a self-contained machine: its own cores, its own RAM, its own attached storage. Nothing is shared between nodes — no shared buffer pool, no shared lock manager, no shared address space, no cross-node memory access. The only way node A can see a row that node B owns is to ask B to send it over the network. This is the opposite of a shared-disk design, where every node mounts the same storage and any node can read any block, and of a shared-memory design (an ordinary multi-core database server), where all threads see one address space. Shared-nothing scales further precisely because there is no shared resource to contend on, but it pays for that with a network hop every time data has to meet data. ## How rows are assigned to nodes An MPP engine needs a deterministic rule that maps each row to exactly one node. Three rules dominate: - **Hash placement.** A column (or set of columns) is chosen as the distribution key; the engine computes a hash of the value and takes it modulo the number of placement buckets. All rows with `customer_id = 42` land on the same node. This is what makes co-located joins possible. - **Round-robin or random placement.** Rows go to nodes in rotation. Storage is perfectly even and loading is fast, but no key is co-located, so any join or grouping will need a network redistribution. - **Replication.** A small table is copied in full to every node. Storage is multiplied by the node count, but joins against it never need to move anything. Engines differ in vocabulary and in how much of this is automatic — some ask you to declare the key in DDL, some choose and re-choose it for you, and cloud engines that keep files on object storage assign disjoint *file sets* to workers rather than pinning rows to a machine forever. The execution consequence is the same in all of them. ## Slices, workers and the real unit of parallelism A node does not process its share single-threaded. It splits it again into parallel units — commonly called slices, workers, or threads, one per core or per storage volume. If a 10-node cluster runs 8 units per node, the table is effectively split 80 ways and 80 scan streams run at once. This matters for two reasons. First, effective parallelism is the *unit* count, not the node count. Second, skew is measured at the unit level: if one unit holds ten times the rows of its peers, the query waits for that one unit even though 79 others finished. ## One plan, many copies The optimizer compiles a query once. The resulting plan is cut into fragments, and every worker executes an identical copy of a fragment over its own local rows. Where a fragment needs rows that live elsewhere, the plan contains an explicit data-movement operator (an exchange or shuffle). Everything below that operator is fully local and embarrassingly parallel; everything above it depends on the network. A coordinator (or leader) node parses the SQL, plans it, ships fragments to the workers, and collects results. In most designs it holds no base table data. It does execute the final merge for a global `ORDER BY ... LIMIT`, which is why a query returning millions of rows to the client can be bottlenecked on a single machine even though the scan was perfectly parallel. ## Storage/compute-separated variants Modern cloud warehouses store table files in object storage rather than on node-local disks, and spin compute clusters up and down against them. That changes durability and elasticity, but not the execution model described here: each worker is still handed a disjoint set of files, still scans them independently, still caches them on local SSD, and still has to exchange rows over the network to meet another worker's data. "Shared-nothing" describes how the *execution* is partitioned, which stays true even when the bytes live in a shared object store. ## Why the layout dominates performance Almost every MPP performance question reduces to this layout. If a join's key matches the placement key on both sides, matching rows are already on the same node and the join is local. If it does not, the engine must redistribute — network cost, memory pressure, spill risk. If the placement key has a few dominant values, one node gets a disproportionate share and becomes the whole query's runtime. Understanding where rows live is the prerequisite for reasoning about all of it.

  • What does a shared-nothing cluster's coordinator node actually do during a query?
    It parses and plans the SQL, cuts the plan into fragments, ships them to the compute nodes, and gathers the final result stream. It normally stores no base table data. Because the last merge step for a global ORDER BY or LIMIT happens there, returning very large result sets serializes through one machine even when the scan was fully parallel.
  • If replicating a table to every node removes shuffle entirely, why not replicate everything?
    Storage and load cost multiply by the node count, and every insert or update must be applied on every node. Replication is only sensible for tables that are small relative to node memory and change infrequently. Replicating a large fact table would consume the cluster's storage and make loading unbearably slow.

Think of a library split across ten branches with no inter-branch catalogue lookup. Each branch can only search its own shelves; comparing books across branches means physically shipping them to one place first.

saying these in an interview costs you the question

  • Thinking any node can read any other node's local data
  • Believing the coordinator holds a copy of all the data
  • Assuming parallelism equals node count, ignoring per-node workers
  • Calling round-robin placement good because storage is even
  • Assuming rows are assigned to nodes randomly in every engine

context