Pattern · seen in 2 breakdowns across 2 companies

Shard-Key Colocation

When you shard a database across machines, keep data that is used together on the same machine, by sharding all the related tables by one shared key - like the user or workspace.

The mechanism

At its core: a database can keep a transaction atomic (all-or-nothing) only within one machine. So if related data is spread across shards, a transaction over them can half-fail and corrupt the data. Shard all the related tables by the same key, and each entity's rows stay on one machine, where transactions are safe.

LIVE ARTIFACTONE ENTITY, ONE SHARDOPEN FULL SCREEN ↗
THE IDEAWhen you shard a database, a transaction that must be atomic (all-or-nothing) only stays atomic within a single machine. Spread the rows a transaction touches across different shards and it can partly succeed - a parent row deleted on one machine while its child rows, stranded on another, never get updated. That is silent corruption. Colocation prevents it: shard the whole group of related tables together, and partition them all by one key - a user, a workspace, an account - so every row for a given one of those lands on the same machine. Then the queries and transactions you actually run stay on a single shard, fast and safe, with no cross-machine coordination.
WHAT TO TRYShard the rows by row id and run a transaction across a workspace - it lands on several machines, one fails, and you get orphaned rows. Colocate by workspace instead and the same transaction stays on one machine and commits cleanly.

Scatter a workspace's rows across shards and a delete can strand its comments; colocate them and the transaction stays on one machine.

Definition

Sharding splits a database across many machines. The catch is that a database can only promise a transaction - a set of changes that all succeed or all fail together - within a single machine. So if two rows that must change together end up on different machines, a change can half-succeed: it commits on one machine and fails on the other, leaving the data inconsistent. Shard-key colocation avoids that by keeping data that is used together on the same machine.

Getting there takes two matching choices:

  • shard the tables together - shard the whole set of tables that reference each other, not just the one big table
  • split them by one key - partition them all by the same key, the one whose natural boundary matches how the application reads and changes data

Pick the user, the workspace, or the account, and every row belonging to one of them lands on one machine. When your own code controls the sharding - see application-layer sharding - this sharding key is the main thing you are choosing.

KEEP RELATED ROWS TOGETHER
Left, a workspace's rows scattered across shards, a transaction can half-succeed; right, all its rows on one shard.
Shard related rows by different keys and a workspace's parent and children land on different shards, so a transaction across them can half-succeed; partition them all by the same key and they stay on one shard, where transactions are safe.

Do this and the queries and transactions you actually run stay inside one machine, where they are fast and safe - no coordinating a change across machines, no pulling a query's data from several of them. Get it wrong and you are back to those cross-machine operations for everyday work, or worse, the silent corruption of related data drifting out of sync. The real work is choosing that key well: real queries should stay inside one entity, and no single entity should be so large it overloads its machine.

When it applies

01You are sharding a relational database whose tables link to each other. When related tables (joined by foreign keys) take part in the same transactions, they need to land on the same machine or those transactions break.
02You are picking a shard key for a multi-tenant product. Usually the tenant - the customer or workspace - is exactly the boundary your transactions and queries stay within, which makes it the natural key to shard by.
03Without colocation, your everyday queries would have to span machines. In any sharded system where the normal path would otherwise mean joining across machines or coordinating a commit across them, colocating removes the scatter-and-gather cost.

Tradeoffs

The shards inherit the entity's uneven sizes: splitting by count looks even, but if one tenant is bigger than a thousand others, its machine carries far more load than the rest.
The key you pick is very hard to change later. It is baked into the whole layout, so if the product shifts and the old key no longer matches how people use it, colocation stops helping and starts fighting you.
Everything connected has to shard together, even tiny tables. A small lookup table that a big table's transactions touch has to be sharded by the same key too, so the scheme pulls in more tables than you might expect.

The same move, 2 ways

Every row is a production system that bet on this pattern — the note says how, in that system's own terms.

Figma
Figma Blog
2024
Figma's 'colos' are exactly this pattern: related tables grouped so they share one shard key and one physical layout. Inside a colo, cross-table joins and full transactions work, as long as they stay within a single value of the shard key. Most application code already queried that way, which is what made the abstraction cheap for product developers: the payoff is that the queries you were already running stay free. Read the breakdown →
Notion
Notion Blog
2021
Both halves of the pattern sit in one post. Shard every table reachable from block so that rows which must change together live on the same host (a single database is as far as a commit-together guarantee reaches), and partition all of it by workspace ID so that a given user's queries land on one shard. The example the post gives is the failure case stated plainly: a block deleted on one host while its comments, stranded on another, never get the update. Read the breakdown →

Often used together

Patterns sharing breakdowns with this one — derived from co-occurrence, threshold ≥2 shared.

Problems this pattern answers

The walls where its breakdowns live — each opens the cross-company comparison.