Sharding Operations Lab
Can a shard plan stay balanced when one key gets hot, an owner fails, or data must move online?
Sharding assigns each key to one authoritative storage owner. The placement rule controls locality and balance; replication protects availability but does not split a hot key or remove the cost of moving ownership.
Step 1
Model
Choose hash, range, or directory placement; shape the key distribution and shard count; then change traffic, write share, replication, read routing, and the rebalance budget.
Step 2
Observe
Trace exact key ownership, read and write paths, per-shard load, busiest-owner headroom, availability, replica lag, movement cost, and the consequence users experience.
Step 3
Challenge
Compare a hot key, tenant skew, shard loss, online reshard, replica lag, and recovery catch-up. Find a configuration that stays available without hiding stale reads or exhausting the busiest shard.
Place keys, route traffic, then break a shard
A good shard plan balances ownership and request pressure while keeping failures and online movement inside a measured capacity budget.
Challenge the design
Balanced keys, normal routing, and all replicas caught up.
Place keys
Change ownership and observe skew.
Route and recover
Change pressure, redundancy, and movement.
100.00%
Modeled request success
50%
Shard 1 sets the limit
80% next
High movement with simple modulo hashing
70 ms
Inside the modeled threshold
Live topology
Keys and requests converge on an owner
Hottest shard receives 25% of traffic
Workload
8.4K reads/s
3.6K writes/s
Router
hash(key) mod shard count
+0.3 ms
Shard 1
1M keys · 3K req/s
Shard 2
1M keys · 3K req/s
Shard 3
1M keys · 3K req/s
Shard 4
1M keys · 3K req/s
Faster, possibly stale
Client → router → nearest replica (RF 2)
One owner coordinates the write
Client → router → owner primary → 1 replica
User-visible consequence
Every owner retains operating headroom
The baseline is healthy, but test a hot key and shard loss before treating the design as production-ready.