PrepZone Logo
PrepZone

Replication and Sharding

Master-replica reads, partition keys and the trade-offs of splitting data across nodes.

Replication basics

RDS primary + read replicas

writeasyncasyncreadCOMPUTE
EKS APIread/write split
DATABASE
RDS primaryus-east-1a
DATABASE
Read replica 1dashboards
DATABASE
Read replica 2EU analytics
Writes hit primary; read-heavy queries route to replicas with bounded staleness.

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.

Java
-- 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:

AspectRead from replicaRead from primary
LatencyLower load on primaryHigher primary load
FreshnessMay miss last few writesAlways latest
StreamHub usePublic feeds, browse pagesPost-upload confirmation, billing
  • Latency

    Read from replicaLower load on primary
    Read from primaryHigher primary load
  • Freshness

    Read from replicaMay miss last few writes
    Read from primaryAlways latest
  • StreamHub use

    Read from replicaPublic feeds, browse pages
    Read 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:

Java
# 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

StrategyHow it worksTrade-off
Rangeuser_id 1–1M on shard AHot ranges if IDs cluster
Hashhash(user_id) % NEven spread; range queries cross shards
DirectoryLookup table maps key → shardFlexible; lookup service is SPOF
GeographicEU users on EU shardCompliance; cross-region queries hard
  • Range

    How it worksuser_id 1–1M on shard A
    Trade-offHot ranges if IDs cluster
  • Hash

    How it workshash(user_id) % N
    Trade-offEven spread; range queries cross shards
  • Directory

    How it worksLookup table maps key → shard
    Trade-offFlexible; lookup service is SPOF
  • Geographic

    How it worksEU users on EU shard
    Trade-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.

Java
# 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.
Java
-- 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:

Java
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:

  1. Single primary + vertical scale.
  2. Read replicas for read-heavy paths.
  3. Cache hot keys in Redis.
  4. 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.