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 UPDATE it needs to be atomic => need for distributed transactions

Sharding + replication

  • Shard replicas. Each shard is distributed among N replicas. 1 replica is a leader. Image

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
  • 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 ZooKeeper or etcd.
    • 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 ZK or etcd
    • 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
  • 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/