Replication basics
RDS primary + read replicas
Replication goals
- Read scaling — route SELECT queries to replicas.
- High availability — promote replica if primary fails.
- Geographic proximity — read replica in EU for EU users.
- Backup / DR — delayed replica for point-in-time recovery.
Leader-follower (primary-replica) is the most common pattern: all writes go to the primary; replicas apply the write-ahead log asynchronously or semi-synchronously.
-- App routing (conceptual)
-- writes → primary.db.streamhub.internal
INSERT INTO videos (id, creator_id, title) VALUES ($1, $2, $3);
-- reads → replica.db.streamhub.internal (accept replica lag)
SELECT id, title, view_count FROM videos WHERE creator_id = $2 ORDER BY created_at DESC LIMIT 20;
Replication lag and consistency
Replicas lag the primary by milliseconds to seconds. Implications for StreamHub:
| Aspect | Read from replica | Read from primary |
|---|---|---|
| Latency | Lower load on primary | Higher primary load |
| Freshness | May miss last few writes | Always latest |
| StreamHub use | Public feeds, browse pages | Post-upload confirmation, billing |
Latency
Read from replicaLower load on primaryRead from primaryHigher primary loadFreshness
Read from replicaMay miss last few writesRead from primaryAlways latestStreamHub use
Read from replicaPublic feeds, browse pagesRead from primaryPost-upload confirmation, billing
After a creator uploads a video, StreamHub reads from primary for the "your video is live" page; followers see it on replicas within ~1 second.
Failover
When the primary dies, a replica is promoted:
# Managed RDS failover (automatic)
detection: health_check_failure_60s
action: promote_read_replica
dns_update: primary.endpoint → new_primary
app_impact: brief_write_unavailability (~30-120s)
Apps must retry writes and handle connection resets during failover. Use connection pools with validation queries.
When to shard
Shard when:
- Write QPS exceeds one primary's capacity (~5–15K sustained writes/sec depending on hardware).
- Storage exceeds comfortable single-node limits (multi-TB with heavy churn).
- Blast radius reduction — isolate tenants or regions.
StreamHub began sharding at ~3M DAU when upload spikes during live events saturated the primary.
Sharding strategies
| Strategy | How it works | Trade-off |
|---|---|---|
| Range | user_id 1–1M on shard A | Hot ranges if IDs cluster |
| Hash | hash(user_id) % N | Even spread; range queries cross shards |
| Directory | Lookup table maps key → shard | Flexible; lookup service is SPOF |
| Geographic | EU users on EU shard | Compliance; cross-region queries hard |
Range
How it worksuser_id 1–1M on shard ATrade-offHot ranges if IDs clusterHash
How it workshash(user_id) % NTrade-offEven spread; range queries cross shardsDirectory
How it worksLookup table maps key → shardTrade-offFlexible; lookup service is SPOFGeographic
How it worksEU users on EU shardTrade-offCompliance; cross-region queries hard
StreamHub shards by creator_id — a creator's videos, stats, and comments live on one shard. Fan-out feeds cross shards but are assembled asynchronously.
# Shard routing (pseudocode)
def shard_for_creator(creator_id: str) -> str:
bucket = consistent_hash(creator_id) % NUM_SHARDS
return f"shard-{bucket}.db.streamhub.internal"
Cross-shard challenges
Sharding complications
- Joins — no cross-shard SQL joins; denormalise or aggregate in app layer.
- Transactions — no global ACID; use Saga or per-shard transactions.
- Rebalancing — adding shards requires moving data (consistent hashing helps).
- Unique constraints — global uniqueness needs central service or UUID.
-- Anti-pattern: cross-shard join
-- SELECT * FROM shard_a.videos v JOIN shard_b.users u ON ...
-- Instead: fetch IDs from shard A, batch-fetch from shard B in application
Replication + sharding together
Each shard is its own primary with its own replicas:
Shard 0: Primary₀ → Replica₀ₐ, Replica₀ᵦ
Shard 1: Primary₁ → Replica₁ₐ, Replica₁ᵦ
Shard 2: Primary₂ → Replica₂ₐ, Replica₂ᵦ
Writes scale with shard count; reads scale with replicas per shard.
Migration path
Do not shard on day one. StreamHub's ladder:
- Single primary + vertical scale.
- Read replicas for read-heavy paths.
- Cache hot keys in Redis.
- Shard when write metrics force it — with dual-write or backfill migration.
Quick recall
Everything you need if you only revisit this box.
- Replication copies data for read scale and failover; primary handles all writes.
- Replica lag means stale reads — route freshness-sensitive queries to primary.
- Shard when writes or storage exceed one primary; pick a stable partition key.
- StreamHub shards by creator_id to co-locate creator data.
- Cross-shard joins and global transactions do not work — design around them.
- Each shard can have its own replicas; scale writes and reads independently.
Test yourself
Answer these before moving on — recall is what makes it stick.