ddia/ch07.md - DDIA
Sharding problems
- relational data, it's hard to find shards for secondary indexes.
- Solutions: scatter-gather; maintain a global secondary index;
- updating data on several shards. assume
UPDATEit needs to be atomic => need for distributed transactions
Sharding + replication
- Shard replicas. Each shard is distributed among N replicas. 1 replica is a leader.

How to choose a sharding strategy
- Key range
- Good: effective range scans
- Problem: hot shards in uneven distributions
- Rebalancing:
- In some systems (like HBase), shard split is triggered when shard reaches (default) 10 GB.
- In others, when write throughput being persistently above a certain threshold
- Overall, very expensive
- Note: if you have fixed number of shards => if data grows it's GG => needs rebalance
- Key hash
- Hash modulo number of nodes
- Problem: after adding new node you have to rebalance everything
- Bad approach
- Hash modulo number of shards
- Better: you can reassign shards as you add/remove nodes, solving the problem above.
- Used in Elasticsearch, Riak, Citus
- Problem: hard to make a correct guess on how many shards you are going to use
- (continuous) Hash range
- Idea: apply hash and assign each shard a (continuous) range of values it is responsible for.
- Trick: partition by one key and sort by key by which you query inside
- Used in DynamoDB, MongoDB (configurable)
- (random) Hash range
- Idea: as in prev. case but ranges can be scattered. results in a fair share of the dataset.
- Used by Cassandra and ScyllaDB
- Hash modulo number of nodes
- Consistent hashing
- Motivation: move as little data as possible
- Consistent hashing is a hash function that maps keys to specified number of shards in a way that satisfies 2 properties:
- The number of keys mapped to each shard is roughly equal
- When the number of shards changes, as few keys as possible are moved from one shard to another
- Random hash range approach above IS consistent hashing
- Idea: have a ring of hash range -> key is hashed to a position on a ring -> walks clockwise -> first node to hit is where it must be stored
- Additionally: how to solve imbalance problem - add virtual nodes
Skewed workloads
- Problem: some shards became hot because of hot keys
- Solutions
- (if the workload is read-heavy), replicate the shard
- (if the workload is read-heavy), hide shard behind caching
- Move hot keys to separate shards
- Needs monitoring for high-frequency shards (heuristics for exceeding thresholds).
- When traffic drops => merge back
Request routing
- Problem with a shared routing layer (like LB): how does it know about changes in the assignment of shards to nodes?
- (possible) Solution: rely on distributed coordination service, like
ZooKeeperoretcd.- Used by HBase, Kubernetes
- "Each node registers itself in ZooKeeper, and ZooKeeper maintains the authoritative mapping of shards to nodes. Other actors, such as the routing tier or the sharding-aware client, can subscribe to this information in ZooKeeper. Whenever a shard changes ownership, or a node is added or removed, ZooKeeper notifies the routing tier so that it can keep its routing information up-to-date" ![[Pasted image 20260509212330.png]]
- (alternative) Solution: use gossip protocols
- TODO: not clear how does it determine the node to shard on
How to shard
- By default, AVOID sharding if you don’t need it. Main reason TO SHARD is scale.
- Decide what you want to shard
- what is the highest-priority data type for sharding? take it and them transitively shard all dependent types (reachable by FK) to avoid costly joins and need for distributed transactions.
- don't shard data doesn't need to be sharded or is difficult to shard (global).
- Decide the sharding key
- determined by type of queries you do
- also decide how to sort data on shard. again, based on nature of queries
- Decide number of shards
- Shard number must have many factors
- Have many shards per physical database -> easier shard migration
- How do you route requests? Mapping sharding key to physical database
- fully static: in app mapping using hash ranges
- dynamic: using coordination service, like
ZKoretcd - celebrity problem: route demanding clients separately
Problems of staying single-instance, a.k.a. when to shard
- TXID wraparound
- (potential) solution: aggressive vacuuming
- Vacuuming in Postgres
- (potential) solution: Partition the data on single instance => enables per-partition vacuuming
- Shared resources like locks
- not clear how to solve
- Table bloat due to MVCC
- (potential) solution: aggressive autovacuuming
- Schema migrations. how would you migrate schema?
- not clear how to solve
- Connection exhaustion
- solution: connection pooling, i.e.
PgBouncer
- solution: connection pooling, i.e.
- Write-heavy workloads. Single primary = death
- no solution
- Recovery and backup time - takes days
- no solution
References:
- Sharding strategy at Notion
- https://www.notion.com/blog/sharding-postgres-at-notion
- https://news.ycombinator.com/item?id=28776786
- How to manage hotspots - S3
- https://www.allthingsdistributed.com/2023/07/building-and-operating-a-pretty-big-storage-system.html
- Partitioning GitHub
- https://github.blog/engineering/partitioning-githubs-relational-databases-scale/