skip to content

Distributed relational systems such as Citus and Vitess ask you to declare, per table, whether it is distributed on a key or replicated to every shard. How would you make that declaration across a real schema, and why does the choice of distribution column matter so much for joins and transactions?

level: principalimportance: nice to knowfreq 26%

answer

  1. Distributed / reference / local
  2. Same column, same shard = colocation
  3. Non-colocated join = broadcast or repartition
  4. Distribution column inside every unique key
  5. Denormalise tenant_id onto child tables

basics

~20 s

Distribute the large, tenant-owned tables on the same key so matching rows land on the same shard; replicate small, shared lookup tables to every shard. Colocated tables can be joined and updated in one local transaction; non-colocated ones force cross-shard joins and distributed commits.

solid answer

~60 s

You classify every table into three buckets: - **Distributed tables** — the big, per-tenant fact tables. All of them are distributed on the *same* logical column (`tenant_id`, `account_id`), so rows with the same key value land on the same shard. Systems call this **colocation**, and it is what makes joins, foreign keys and multi-statement transactions across those tables execute entirely locally. - **Reference / replicated tables** — small, slow-changing, shared-by-everyone data (currencies, plans, country codes). A full copy on every shard makes joins to them local everywhere; the price is that writes must be applied to all shards, so they must be rare. - **Local / coordinator-only tables** — metadata that never joins to sharded data. The distribution column must generally be part of the primary key and every unique constraint, because uniqueness is only enforceable within a shard, and foreign keys between distributed tables only work when both sides are colocated on that column. Get this wrong — distribute orders on `order_id` and customers on `customer_id` — and the most common join in your system becomes a repartition or broadcast join at query time.

code

sql · 11 lines
sql
SELECT create_distributed_table('customers',   'tenant_id');
SELECT create_distributed_table('orders',      'tenant_id');
SELECT create_distributed_table('order_items', 'tenant_id');
SELECT create_reference_table('currencies');

CREATE TABLE orders (
  tenant_id  bigint NOT NULL,
  order_id   bigint NOT NULL,
  created_at timestamptz NOT NULL,
  PRIMARY KEY (tenant_id, order_id)
);

go deeper

for a junior

Know that such systems require you to say how each table is spread, and that small lookup tables are copied everywhere.

for a middle

Define colocation and explain that joining tables distributed on the same key stays local while others must broadcast or repartition.

for a senior

Design the whole schema's declaration, including denormalising the distribution column onto child tables and the constraint rules that follow from it.

for a principal

Judge whether the workload has a dominant entity at all, weigh the irreversibility of the distribution column, and be willing to say that a genuinely cross-entity workload should not be forced into one.

## The model these systems impose Citus (a PostgreSQL extension) and Vitess (a MySQL-based system) both take an existing relational database and spread tables over many nodes while keeping SQL. They both require you to declare, per table, how it participates: - a **distributed / sharded table** carries a distribution column (Citus) or a keyspace vindex column (Vitess); rows are placed by hashing (usually) that column into shards; - a **reference / replicated table** is copied in full to every shard; - a **local table** lives only on the coordinator and does not participate in distributed queries. The declaration is not a hint. It determines what the engine can execute locally and what needs network movement. ## Colocation is the central idea Two distributed tables are **colocated** when they use the same distribution column semantics and the same bucket assignment, so for any key value K, the rows of both tables with key K live on the same shard. Why that matters: - **Joins.** `orders JOIN order_items USING (tenant_id, order_id)` where both are distributed on `tenant_id` becomes, on each shard, an ordinary local join over local data. The coordinator simply unions the per-shard results. Without colocation, the engine must either **broadcast** one side to every shard or **repartition** both sides over the network by the join key — orders of magnitude more expensive, and memory-hungry. - **Foreign keys.** A foreign key can only be enforced where both rows are visible. Colocated tables can carry real referential integrity (with the distribution column in the key on both sides); non-colocated ones cannot, and you are left enforcing it in application code. - **Transactions.** Updating several colocated tables for one tenant is a single-node transaction with ordinary ACID semantics and ordinary performance. Touching two shards requires a distributed commit (two-phase commit), which adds latency, introduces in-doubt transactions to operate, and can block on coordinator failure. - **Unique constraints.** A `UNIQUE (email)` on a distributed table cannot be enforced globally, because two conflicting rows may hash to different shards. These systems therefore require unique and primary keys to *include* the distribution column. If you need global uniqueness on something else, you build it — a separate reference table, an external id service, or you accept UUIDs. ## Designing the declaration for a real schema Start from the workload, not the ER diagram. Ask: what value is present on essentially every request? In B2B SaaS that is the tenant; in consumer products often the user. That becomes the distribution column, and you add it to every large table that belongs to that entity — frequently denormalising it onto child tables that previously reached the tenant only transitively through their parent. Yes, this duplicates a column onto `order_items`; that duplication is the price of local joins and is standard practice. Then sort the remainder: - Small and read-mostly, joined from many places → reference table replicated everywhere. Currencies, product catalogues under a few hundred thousand rows, feature flags. Watch the write cost: every insert touches every shard, so a "reference" table with a busy write path is a mistake. - Large but *not* tenant-scoped, joined to tenant data → the hard case. Options are to duplicate it per tenant, to accept a broadcast join, or to reconsider whether the query belongs in the transactional system at all. - Metadata unrelated to sharded data → leave local. ## What you give up, and how to reason about it Some queries will always cross shards: analytics over all tenants, admin searches by non-key columns, global ordering. Both systems can execute those — Citus pushes fragments to workers and merges, Vitess scatters and merges at vtgate — but the cost model is different: latency is bounded by the slowest shard, results may need a full sort at the merge point, and a scatter consumes a connection per shard. It is reasonable to route that class of query to read replicas, a columnar table, or a downstream analytical store rather than paying for it on the OLTP path. You also inherit operational shape: schema changes must be applied to every shard consistently (both systems automate this, but it becomes an online-DDL problem at scale), and shard rebalancing becomes a routine capacity operation. ## The principal-level judgement The real decision is not "which column" but **whether the application's access patterns actually have a dominant entity**. If 95% of queries carry a tenant id, a distributed relational system gives you near-single-node semantics with horizontal capacity, and it is an excellent trade. If your access patterns are genuinely graph-shaped or cross-entity — every query joins across users, or the product is fundamentally about relationships between tenants — no distribution column will save you, and forcing one produces a system where the common query is a repartition join. In that case the honest options are a bigger single node, functional decomposition into services with their own databases, or a datastore designed for that shape. Also weigh reversibility: distribution column choice is the least reversible decision in the schema. Prefer a key you already know is present in the API surface, keep placement indirect so keys can be moved, and prototype the top ten queries against a two-shard cluster before committing — a query that scatters at two shards will scatter at fifty. ## How to answer Classify tables into distributed / reference / local, define colocation and connect it to joins, foreign keys, transactions and unique constraints, show that you would denormalise the distribution column onto child tables, and close with the judgement call: distribution works when a dominant entity exists in the access pattern, and no key choice rescues a workload where one does not.

  • Why do these systems insist that the distribution column be part of the primary key and every unique constraint?
    Because uniqueness can only be checked within one shard without a synchronous cross-shard round trip. If the constraint does not include the distribution column, two conflicting rows can hash to different shards and neither shard can see the violation. Including the column guarantees that all rows sharing the constrained values are on the same node, making the check a normal local index probe.
  • You have a large table that is not scoped to a tenant but is joined to tenant data constantly. What are your options?
    Three, none free. Replicate it to every shard if it is small enough and rarely written, accepting the write amplification. Denormalise the needed columns into the tenant-scoped tables, accepting duplication and an update path. Or accept a broadcast/repartition join and confine it to queries that tolerate the latency, ideally on replicas. If none is acceptable, that table is a signal that the workload may not fit a single distribution column.

Colocation is keeping one customer's whole case file in one office. Reference tables are the price list pinned to the wall of every office.

saying these in an interview costs you the question

  • Distributing parent and child tables on different columns and expecting joins to stay local
  • Believing a global UNIQUE constraint on a non-distribution column is enforced across shards
  • Making a write-heavy table a reference/replicated table and ignoring that every write hits every shard
  • Assuming cross-shard transactions behave like local ones rather than requiring two-phase commit
  • Treating the distribution column as a setting that can be changed later without a full data migration

context