A single relational database server can no longer absorb your write volume, so the team proposes splitting the rows of the biggest tables across many independent database servers. What is sharding, and what properties make a column a good shard key?
answer
- Independent servers, disjoint rows
- Key must be in the predicate
- Cardinality caps shard count
- Even by traffic, not just rows
- Colocate the tenant's tables
basics
~20 sSharding splits the rows of one logical table across several independent database servers, choosing the server by a shard key. A good shard key appears in nearly every query, has high cardinality, spreads load evenly, and keeps rows that are read together on the same shard.
solid answer
~60 sSharding is horizontal partitioning across **separate database instances**: each shard is its own server with its own storage, buffer pool, WAL and connection pool, holding a disjoint subset of rows. A routing function maps a row's **shard key** to a shard, so writes and reads scale by adding machines rather than growing one. A good shard key has four properties: - **Ubiquity** — it appears in almost every query's predicate, so requests hit one shard instead of fanning out to all of them. - **High cardinality** — enough distinct values that you can keep splitting; a boolean or a status column caps you at a handful of buckets. - **Even distribution** — of both rows *and* traffic; a key that is uniform by row count can still be hot by request rate. - **Colocation** — related rows (a tenant's orders, invoices and users) land on the same shard, so joins and transactions stay local. In practice the key is usually a tenant/customer id or a user id, not an auto-increment surrogate id, because it is what the application already knows on every request.
go deeper
Be able to say sharding splits rows across separate database servers and that the shard key decides which server, and name one sensible key like customer id or user id.
Name the four properties of a good key and connect a wrong key to concrete pain — scatter-gather reads and a hot shard.
Argue the choice for a specific workload, cover colocation of related tables, secondary-lookup access paths, and the operational cost of N independent databases.
Frame it as a last-resort architectural commitment: state the alternatives you would exhaust first, the blast radius of a wrong key, the migration path if the key must change, and what consistency guarantees you are willing to give up.
## What sharding actually is Sharding is horizontal partitioning where the pieces live on **different database servers**. Each shard is a full, independent database: its own process, memory, disk, write-ahead log, locks and connection pool. Shard 3 does not know shard 7 exists. Together they hold one logical table, split by rows — every row belongs to exactly one shard. That independence is the whole point and the whole cost. It is why sharding scales writes (each shard commits its own transactions on its own disk, so aggregate write throughput grows roughly linearly with shard count) and why it hurts (there is no global transaction, no global unique index, no global ORDER BY, unless something above the shards builds one). This is different from partitioning a table *inside* one server, which splits storage but shares the same CPU, memory and WAL — that buys manageability and plan pruning, not write throughput. ## The shard key and the routing function Every request must answer one question before it can touch data: *which server holds this row?* The answer comes from the **shard key** — one column (or a small composite) present in the row — fed through a routing function: - hashing the key and taking the result modulo the shard count, or mapping the hash into a fixed ring of virtual buckets; - comparing the key against range boundaries; - looking the key up in a directory table that stores an explicit key→shard mapping. The shard key is chosen once, early, and is extremely expensive to change later, because changing it means physically moving essentially every row. Teams routinely spend more design time on this one column than on the rest of the schema. ## What makes a key good **It is in the query.** If the application asks "give me order 4711" but the table is sharded by `customer_id`, the router has no idea where order 4711 lives and must ask every shard — a *scatter-gather*. Scatter-gather costs the latency of the slowest shard, multiplies connection usage, and degrades as you add shards. Sharding by the value the application already carries on every request (usually the tenant or the user) turns almost all traffic into single-shard traffic. A common fix for secondary lookups is to embed the shard key in the visible identifier, so an order id literally encodes its customer's shard. **It has high cardinality.** The number of distinct key values is a hard ceiling on how far you can split. `country_code`, `status`, `is_active` and `region` all look tempting and all cap out: with five statuses you can never have more than five shards, and one of them will be gigantic. **It distributes evenly — by traffic, not just by rows.** Row-count balance is easy to check and easy to be fooled by. A B2B system sharded by `tenant_id` may have a perfectly reasonable row distribution and still put 40% of all queries on the shard holding your largest customer. **It colocates what belongs together.** If users, orders, order lines and invoices all carry the same `tenant_id` and are all sharded on it, then a tenant's entire working set lives on one server: joins are local, foreign keys are enforceable, and a multi-table write is an ordinary local transaction rather than a two-phase commit. Losing colocation is what turns a sharded system into a distributed-systems project. ## Typical choices and their consequences - **Tenant / customer id** — the default for B2B SaaS. Great colocation and ubiquity; risk is tenant skew, since real customer sizes follow a power law. - **User id** — the default for consumer systems. Cardinality is huge and per-user traffic is fairly bounded, so distribution is good; cross-user features (feeds, social graphs, leaderboards) become scatter-gather. - **Monotonic surrogate id or timestamp** — fine under a *hash*, poisonous under a *range*, where all new writes land on the newest shard. - **Composite key** — e.g. hash of `(tenant_id)` for placement while the primary key stays `(tenant_id, id)`. Most distributed relational systems require the shard key to be part of the primary key and of every unique constraint, precisely because uniqueness can only be enforced within a shard. ## What sharding costs you Be honest about this in an interview: cross-shard joins and aggregates need a coordinator or an application-side merge; global uniqueness and global sequences need an external source (UUIDs, snowflake ids, a ticket server); transactions spanning shards need 2PC or a saga; every shard needs its own backup, monitoring, failover and schema migration, applied consistently. Sharding is the tool you reach for after read replicas, caching, partitioning, and a bigger box have run out — not before. ## How to answer Define sharding as horizontal split across independent servers, name the four key properties (ubiquitous, high-cardinality, evenly distributed by traffic, colocating), give a concrete key for a concrete workload, and volunteer the price: scatter-gather, no global uniqueness, no free cross-shard transactions.
- The application usually queries by customer_id, but a support tool looks rows up by order_id. How do you serve that lookup without scattering to every shard?Either embed the shard key inside the order id so it is derivable at routing time, or maintain a secondary lookup table mapping order_id to its shard (or to the customer_id) that the router consults first. The lookup table is a second write on the insert path and must be kept consistent, but it turns an N-shard fan-out into two point reads. Scatter-gather is acceptable only for rare, low-volume access paths.
- Why do distributed relational systems usually require the shard key to be part of every unique constraint and primary key?Uniqueness can only be checked cheaply where all candidate rows live. If the constraint does not include the shard key, two conflicting rows can land on different shards and neither shard sees the conflict, so enforcing it would need a synchronous check across all shards on every insert. Including the shard key makes the constraint locally enforceable, which is why global uniqueness on other columns is typically delegated to UUIDs or an external id service.
- Is an auto-increment primary key a good shard key?Usually not. It is high cardinality, which is good, but the application rarely knows it before an insert, it carries no colocation meaning, and under range sharding it is monotonic, so all inserts hit the newest shard. Under hashing it distributes fine but destroys locality, scattering a single tenant's rows across every shard.
Shards are separate filing cabinets in separate offices. The shard key is the rule that tells the clerk which office to walk to. A bad rule means walking to every office and merging what you find.
saying these in an interview costs you the question
- Calling partitioning within one server 'sharding' and claiming it scales writes
- Choosing a low-cardinality column such as region, country or status as the shard key
- Assuming even row counts mean even load, ignoring per-tenant request skew
- Expecting cross-shard joins, global unique indexes or ACID transactions to keep working transparently
- Reaching for sharding before replicas, caching, partitioning and vertical scaling are exhausted