skip to content

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

level: seniorimportance: nice to knowfreq 36%

answer

  1. something must define what a table is
  2. object stores have no multi-object transactions
  3. a commit swaps a list, not the data
  4. one place both clusters ask
  5. statistics for pruning live there too

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.

solid answer

~50 s

Object storage gives you durable bytes and no transactions; a stateless cluster forgets everything when it stops. The consistency lives in a third tier — a metadata, catalog and transaction service, usually backed by a transactional store of its own. It holds the file list that defines each table version, per-file statistics used to skip files during pruning, and the commit log that serializes writers. A write does not modify existing files: it uploads new ones and then commits a new file list in a single metadata operation, so readers atomically see the old set or the new one. That is also what makes independent clusters mutually consistent — they are not gossiping with each other, they are each resolving the same table version through the same service. The trade-off is that this tier is now a shared dependency on the critical path of every query start and every commit.

go deeper

for a junior

Recall that there is a catalog tier separate from both the files and the query machines, and that it is what tells the engine which files a table currently contains.

for a middle

Explain how immutability plus a single atomic metadata commit gives multi-file writes atomicity, and what else the tier stores: statistics, schema, versions.

for a senior

Reason about it operationally — latency on every query start, commit contention from high-frequency writes, and file count rather than data volume as the cost driver.

for a principal

Treat it as the platform's real coordination point and single shared dependency: plan ingestion cadence, retention windows and file consolidation around what that tier can absorb.

## The missing tier Describe a separated architecture as just "data in object storage, stateless compute" and something is obviously unaccounted for. Object stores are key-value blob stores: they give you durability and per-object atomicity and nothing resembling a multi-object transaction. Compute clusters are disposable. So what makes a table a table, and what makes a multi-file write atomic? The answer is a third tier, variously called the metadata layer, catalog, or cloud services layer, and it is stateful and transactional. It typically runs on a conventional transactional store — a distributed key-value system or a relational database — because it needs ACID semantics over small, hot records. ## What it holds - **The file list per table version.** A table is not a directory; it is the set of file identifiers that the metadata layer says currently constitute it. Older sets are retained for a while, which is exactly what enables history retention and querying a table as of an earlier point. - **Per-file statistics.** Minimum and maximum values per column per file, null counts, row counts, distinct estimates. These are what let a query eliminate files without opening them — the pruning path that makes remote storage tolerable. - **The commit log and concurrency control.** Writers are serialized here. A commit validates that the version it started from is still current and then publishes a new version. - **Schema, ownership and access rules**, plus query history and usage accounting. ## Why a write is atomic without transactional storage Because the files are immutable, a write never mutates anything a reader is looking at. An `INSERT` or `MERGE` uploads new data files — possibly many — and, once every upload has succeeded, performs one metadata commit that swaps the table's file set. Readers that resolved the table before the commit read the old list to completion; readers after it see the new one. If the cluster dies halfway through uploading, the orphaned files were never referenced by any committed version and are simply garbage to be cleaned up. This is the whole trick: **multi-file atomicity is bought with single-record atomicity in the metadata store.** It also explains a family of behaviours that otherwise look magical. Creating a copy of a huge table can be instant, because copying is writing a new metadata entry that references the same files rather than moving data. Reading a table as it was an hour ago is resolving an older file list. Rolling back a bad load is republishing the previous version. ## Why it matters operationally **It is the consistency point for every cluster.** Two clusters of different sizes running different workloads agree on what a table contains because both ask the same service, not because they coordinate. Remove that tier and you would need consensus among compute clusters, which is precisely what nobody wants to build. **It is on the critical path.** Every query start resolves tables, versions, statistics and permissions through this tier before any scan begins. A slow metadata tier shows up as latency on trivial queries — the classic symptom of "even a one-row select takes a second". **It is sensitive to file count, not data volume.** A hundred terabytes in large files is a small metadata problem; the same data in tens of millions of tiny files is a large one, because planning must consider each file's statistics. This is why high-frequency small writes — a streaming pipeline committing every few seconds — stress the metadata layer far more than a nightly bulk load of the same volume, and why every such platform pairs ingestion with background file consolidation. **Its retention settings cost storage.** Keeping old versions means keeping the files those versions reference, so a generous history window on a churn-heavy table can hold far more bytes than the table's current size. That is a metadata policy with a storage bill attached. ## The open variant Open table formats implement the same idea with the metadata written into files in the object store alongside the data, plus a catalog to name the current version. The division of labour is identical — an immutable data layer, a versioned list that defines the table, and an atomic pointer swap — even though the mechanism differs, and the details of those formats are their own subject. ## Saying it well in an interview The crisp version: separation of storage and compute is really a *three*-way split, and the third piece is where the database-ness lives. Object storage supplies durability, compute supplies CPU, and the metadata service supplies transactions, versions and the statistics that make remote reads survivable. Candidates who describe only two tiers usually cannot then explain how a multi-file write is atomic, or why two independent clusters never disagree.

  • Why is the metadata tier a scaling concern for streaming ingestion?
    Cost there scales with file and commit count, not bytes. A pipeline committing every few seconds produces enormous numbers of small files, each with statistics that planning must consider and each a separate commit to serialize. Query planning slows and commit contention rises, which is why these platforms pair streaming ingestion with background consolidation of small files.
  • How does this design make copying a large table nearly instant?
    The copy is a metadata operation: a new table entry referencing the same immutable files. No bytes move. Storage is only consumed later, as either copy diverges and writes new files, since the shared files stay referenced by both.
  • What does a slow metadata tier look like from the outside?
    Constant latency added to every statement regardless of size — trivial queries taking as long as small scans, planning times dominating short queries, and commits queueing behind each other. It is distinguishable from a cold cache because the penalty does not shrink when the same query is repeated.

saying these in an interview costs you the question

  • Claims object storage provides multi-object transactions
  • Says clusters coordinate with each other to stay consistent
  • Thinks commits rewrite existing data files in place
  • Describes only two tiers and cannot explain write atomicity
  • Assumes metadata cost scales with data volume rather than file count

context