Topic 9.4
Resharding, Cross-Shard Queries and Global Indexes
In one line
Once sharded, the hard parts are queries that don't include the shard key (scatter-gather or a global secondary index), transactions spanning shards (avoid, or use sagas or 2PC), and resharding live data (copy, catch up via change streams, switch the routing mapping, then clean up).
Think of it like this
Moving a family from one apartment to another while they keep living there. You copy belongings over, keep forwarding new mail, and at one quiet moment swap the keys, then clear the old flat.
Key ideas
- 01
Scatter-gather: send the query to all shards and merge (sort, limit, aggregate). Latency is the slowest shard; load multiplies by the shard count. Acceptable for rare admin queries, not for hot paths.
- 02
Global secondary index: a separate table sharded by the secondary key mapping to the primary shard key (
email → user_id). It's updated asynchronously (eventually consistent) or transactionally with 2PC (expensive). DynamoDB GSIs are asynchronous. - 03
Cross-shard transactions: design so transactional units share a shard key. When unavoidable (transfer between users on different shards), use a saga with a ledger and idempotent steps, or a database with built-in distributed transactions (Spanner, CockroachDB), accepting the latency.
- 04
Resharding a logical shard: (1) snapshot-copy it to the target, (2) stream changes (logical replication or CDC) until caught up, (3) briefly block writes for that shard, verify the offset, switch the mapping, (4) unblock, (5) delete the source copy later. Vitess and Citus automate variants of this.
- 05
Unique constraints across shards (a global unique email) need a global index or a dedicated uniqueness service; per-shard UNIQUE isn't enough.
Code & diagrams
Interview problem
The problem
Look up users by email in a sharded user store
Users are sharded by user_id across 32 shards. Login needs lookup by email, and emails must be globally unique. Design it.
When it breaks
Hot-path scatter-gather
What you see
A search endpoint without the shard key fans out to 64 shards; p99 equals the slowest shard, and load on every shard rises with every request.
Fix & prevent
Add a global index or search index (Elasticsearch via CDC) for that access pattern; reserve scatter-gather for rare admin use.
Explain it without notes
How do you enforce a globally unique field in a sharded database?
Practice
List the steps to split a hot logical shard onto a new node without downtime.
Trade-offs
- ↔
Global indexes enable non-key lookups at the cost of write complexity and eventual consistency; scatter-gather is simple but doesn't scale.
Done when you can
I can design global indexes, avoid cross-shard transactions, and describe live resharding.