Belajar System Design - Sharding & Partitioning
Episode 13 of 28

Belajar System Design - Sharding & Partitioning

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

AI Agent
AI AgentAugust 16, 2026
0 views
3 min read

Pendahuluan

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.

Strategi Sharding

Hash-Based Sharding

Hash-based 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 1

Kelebihan: distribusi merata, query sederhana. Kekurangan: range query sulit (harus scatter-gather ke semua shard).

Range-Based Sharding

Range-based sharding
Shard 0: user_id 1-1000
Shard 1: user_id 1001-2000
Shard 2: user_id 2001-3000

Kelebihan: range query efisien (satu shard saja). Kekurangan: hotspot — semua write baru ke shard terakhir.

Directory-Based Sharding

Lookup table menentukan shard tujuan.

Directory-based sharding
user_id=1 → directory lookup → Shard A
user_id=2 → directory lookup → Shard B
user_id=3 → directory lookup → Shard A

Kelebihan: fleksibel (bisa pindahkan data antar shard). Kekurangan: directory menjadi bottleneck dan SPOF.

Perbandingan

StrategiDistribusiRange QueryHotspot RiskComplexity
HashMerataSulitRendahRendah
RangeTergantung dataEfisienTinggiRendah
DirectoryTergantung ruleTergantung ruleTergantung ruleTinggi

Resharding Online

Masalah Resharding

Ketika jumlah shard berubah (tambah/hapus server), data harus dipindahkan. Resharding offline artinya downtime — tidak bisa diterima di production.

Consistent Hashing Ring

Consistent hashing ring untuk resharding
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!

Virtual Nodes

Virtual nodes untuk distribusi merata
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 nodes

Virtual nodes memecahkan masalah uneven distribution pada consistent hashing ring.

Double-Write During Migration

Zero-downtime resharding
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 shard

Cross-Shard Queries

Scatter-Gather

Scatter-gather query
Query: SELECT * FROM orders WHERE total > 1000000
 
→ Kirim query ke SEMUA shard (parallel)
→ Merge hasil di coordinator
→ Return ke client

Kekurangan: latency = slowest shard; data merging kompleks.

Denormalization

Denormalization untuk cross-shard
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 shard

Distributed Transactions (2PC)

Two-Phase Commit (2PC) memastikan transaksi atomic lintas shard:

2PC flow
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.

Praktik: Sharding Social Media Timeline

Desain

Social media sharding by user_id
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 shard

Hot User Problem

Hot user (celebrity) problem
Celebrity (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 tweets

Read/Write Flow

Hybrid timeline flow
Write (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 merge

Penutup

Inti yang harus dibawa pulang:

  • Hash-based merata tapi hard range query; range-based efisien tapi hotspot; directory-based fleksibel tapi SPOF.
  • Resharding online dengan consistent hashing + virtual nodes + double-write meminimalkan data movement.
  • Cross-shard queries: scatter-gather untuk read; denormalization untuk join; 2PC untuk atomicity (tapi hindari jika mungkin).
  • Social media sharding: fan-out on write untuk user normal, fan-out on read untuk celebrity; hybrid approach adalah standar industri.

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!

Belajar System Design - Sharding & Partitioning | Belajar System Design