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.
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.
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
Tradeoffs
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.
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.