Memahami strategi sharding (hash-based, range-based, directory-based), resharding online dengan consistent hashing ring, cross-shard queries, distributed transactions (2PC), dan desain sharding untuk social media timeline

Setelah di episode 12 kita memahami replication strategies, pada episode ini kita masuk ke teknik horizontal scaling database yang paling powerful: sharding. Di episode 4 kita sudah membahas dasar sharding; di episode ini kita bedah lebih dalam: strategi advanced, resharding, cross-shard queries, dan case study nyata.
Sharding adalah jawaban untuk pertanyaan: "bagaimana jika satu database server tidak cukup untuk menampung semua data?" Ketika satu server sudah mencapai batas storage atau throughput, satu-satunya jalan adalah membagi data ke banyak server — dan itu adalah sharding.
shard = hash(user_id) % number_of_shards
user_id=1 → hash(1) % 4 = 1 → Shard 1
user_id=2 → hash(2) % 4 = 3 → Shard 3
user_id=5 → hash(5) % 4 = 1 → Shard 1Kelebihan: distribusi merata, query sederhana. Kekurangan: range query sulit (harus scatter-gather ke semua shard).
Shard 0: user_id 1-1000
Shard 1: user_id 1001-2000
Shard 2: user_id 2001-3000Kelebihan: range query efisien (satu shard saja). Kekurangan: hotspot — semua write baru ke shard terakhir.
Lookup table menentukan shard tujuan.
user_id=1 → directory lookup → Shard A
user_id=2 → directory lookup → Shard B
user_id=3 → directory lookup → Shard AKelebihan: fleksibel (bisa pindahkan data antar shard). Kekurangan: directory menjadi bottleneck dan SPOF.
| Strategi | Distribusi | Range Query | Hotspot Risk | Complexity |
|---|---|---|---|---|
| Hash | Merata | Sulit | Rendah | Rendah |
| Range | Tergantung data | Efisien | Tinggi | Rendah |
| Directory | Tergantung rule | Tergantung rule | Tergantung rule | Tinggi |
Ketika jumlah shard berubah (tambah/hapus server), data harus dipindahkan. Resharding offline artinya downtime — tidak bisa diterima di production.
Sebelum (3 node): A(0°) B(120°) C(240°)
Sesudah (4 node): A(0°) B(90°) C(180°) D(270°)
Hanya data antara A-B yang perlu dipindah ke B (sebagian)
Hanya data antara B-C yang perlu dipindah ke C (sebagian)
= Minimal data yang dipindah!Physical Node A: VNode A-1 (10°), A-2 (130°), A-3 (250°)
Physical Node B: VNode B-1 (40°), B-2 (160°), B-3 (280°)
Total 6 vnodes → distribusi lebih merata dari 2 physical nodesVirtual nodes memecahkan masalah uneven distribution pada consistent hashing ring.
Phase 1: Start dual-write (write ke shard lama DAN baru)
Phase 2: Backfill existing data dari lama ke baru
Phase 3: Read dari shard baru (verify consistency)
Phase 4: Stop write ke shard lama
Phase 5: Remove old shardQuery: SELECT * FROM orders WHERE total > 1000000
→ Kirim query ke SEMUA shard (parallel)
→ Merge hasil di coordinator
→ Return ke clientKekurangan: latency = slowest shard; data merging kompleks.
Orders shard: user_id → orders
Users shard: user_id → user profile
Untuk query "orders + user name":
→ Join di aplikasi, bukan di database
→ atau denormalize: simpan user_name di orders shardTwo-Phase Commit (2PC) memastikan transaksi atomic lintas shard:
Phase 1 (Prepare):
Coordinator → semua participants: "Can you commit?"
Participants: "Yes, I can" (prepare selesai, lock dipegang)
Phase 2 (Commit):
Coordinator → semua participants: "Commit!"
atau: Coordinator → semua participants: "Abort!"Masalah 2PC: blocking (jika coordinator crash, participants terkunci), latency tinggi (tunggu semua participant).
Warning
2PC sebisa mungkin dihindari karena blocking nature-nya. Jika memungkinkan, desain aplikasi agar tidak butuh cross-shard transactions — gunakan eventual consistency atau saga pattern sebagai pengganti.
Shard key: user_id (hash-based)
Tweets table: user_id, tweet_id, content, created_at
Timeline table: user_id, tweet_id (denormalized)
Write: user posts tweet → insert ke tweets shard → fan-out ke followers' timeline
Read: user scroll timeline → read dari自己的 timeline shardCelebrity (10M followers): post 1 tweet
→ Fan-out ke 10M timeline entries (sangat mahal!)
Solusi:
1. Hybrid approach:
- User dengan <10K followers: fan-out on write (push)
- Celebrity (>10K followers): fan-out on read (pull)
2. Cache celebrity's recent tweets
3. Timeline = merge of pushed + pulled tweetsWrite (normal user):
1. Save tweet
2. Push ke timeline semua followers (<10K)
Write (celebrity):
1. Save tweet
2. TIDAK push (terlalu banyak)
3. Cache recent tweets
Read:
1. Load pushed tweets dari自己的 timeline shard
2. Pull recent tweets dari celebrities yang di-follow
3. Merge dan sort by timestamp
4. Cache hasil mergeInti yang harus dibawa pulang:
Di episode 14 selanjutnya kita akan membahas fault tolerance & resilience patterns — circuit breaker, bulkhead isolation, retry with exponential backoff, graceful degradation, dan chaos engineering. Sistem yang tidak bisa gagal dengan selamat bukanlah sistem yang baik!