🏫 The School›📐 System Design›⚖️ धडा 06 — Load balancing + consistent hashing: पाहुणे आणि keys पसरवा
🖼️ See the drawing + lab 🏠 Course home 🌿 Branch on GitHub ✏️ View source
🖼️ आकृती आणि labThe drawing + lab पूर्ण पानावर उघडा ↗Open full page ↗

⚖️ धडा 06 — Load balancing + consistent hashing: पाहुणे आणि keys पसरवा

📍 तुम्ही इथे आहात: 18 पैकी धडा 06 · मागे: lesson-05-caching · पुढे: lesson-07-queues


📦 या ब्रँचमध्ये काय आहे

धडे 01–05, आणि पसरवण्याचे दोन प्रकार. Load balancers एकसारख्या API servers वर requests पसरवतात: L4 vs L7, algorithms, health checks. Consistent hashing cache (किंवा data) servers वर keys पसरवते, म्हणजे एक server जोडला किंवा गमावला तरी फक्त त्याचाच वाटा हलतो. design/blocks.py मधील HashRing आणि moved_keys() हलणाऱ्या keys मोजतात; design/demo.py मधील balance() hash % N ची ring शी तुलना करते.

🧒 5 वर्षांच्या मुलाला समजावल्यासारखे

फाटकावरची मदतनीस. गर्दीच्या सकाळी 12 कारकून 12 खिडक्यांवर बसतात. फाटकावरची एक मदतनीस प्रत्येक पालकाला एका खिडकीकडे पाठवते: क्रमाने पुढची, किंवा सर्वात छोटी रांग असलेली. दर काही मिनिटांनी ती प्रत्येक खिडकी तपासते — "चालू आहे ना?" — आणि घरी गेलेल्या कारकुनाकडे पालक पाठवणे थांबवते. ती मदतनीस म्हणजे load balancer.

कप्पे. आता वेगळे काम. कार्यालय वर्गांच्या सूचनांच्या प्रती 4 कपाटांत ठेवते. वर्ग 3A कोणत्या कपाटात? दीपिकाचा पहिला नियम: "वर्ग क्रमांक भागिले 4 — बाकी म्हणजे कपाट." तो चालतो. मग 5 वे कपाट येते. आता नियम "भागिले 5" होतो, आणि जवळजवळ प्रत्येक वर्ग वेगळ्या कपाटात जातो. कारकुनांना दर 10 पैकी 8 फाइली हलवाव्या लागतात. सगळ्या प्रती एकदम "हरवतात".

दीपिकाचा अधिक चांगला नियम: भिंतीवर एक मोठे घड्याळ ⏰ रंगवा. प्रत्येक कपाटाला घड्याळावर काही खुणा मिळतात. प्रत्येक वर्गालाही घड्याळावर एक जागा मिळते, आणि तो घड्याळाच्या दिशेने पुढच्या कपाट-खुणेत राहतो. 5 वे कपाट आले की ते आपल्या खुणा जोडते — आणि फक्त आपल्या खुणांच्या लगेच आधीचे वर्ग घेते. सुमारे 5 पैकी 1 फाइल हलते. बाकीच्या आहेत तिथेच राहतात.

🗺️ आकृती

flowchart LR
    subgraph lb["⚖️ load balancer (requests)"]
      c["👪 requests"] --> alb["L7 ALB<br/>round robin / least conn<br/>health checks"]
      alb --> a1["API 1"]
      alb --> a2["API 2"]
      alb --> a3["API …12"]
    end
    subgraph ring["⏰ consistent hashing (keys)"]
      k["10,000 class keys"] --> m["hash % N, 4 → 5<br/>8,001 move (80%)"]
      k --> r["hash ring, 4 → 5<br/>1,975 move (20%)<br/>100 virtual nodes each"]
    end
    a1 --> r

🗺️ काढलेली आवृत्ती + एक lab: https://school-edh.pages.dev/system-design/lesson-diagrams.html#l06

❓ काय

🤔 का

कारण धडा 02 ला 12 API servers लागतात आणि cache 4 nodes वरून 5 वर वाढू शकतो. Load balancer 12 servers ना एकासारखे दाखवतो आणि बिघडलेला लपवतो. Consistent hashing cache ला कड्यावरून न पडता वाढू देते: % N सह, गर्दीच्या सकाळी एक node जोडला की 80% cache रिकामा होतो आणि ते reads अगदी वाईट क्षणी database कडे जातात.

🔧 कसे (या repo मध्ये)

design/blocks.py मधील HashRing(servers, vnodes=100) प्रत्येक server चे vnodes बिंदू ring वर ठेवते ("s0#17" चे MD5 …), bisect ने क्रमवार ठेवलेले. owner(key) key च्या hash नंतरचा पहिला बिंदू शोधते, शेवटी पुन्हा सुरुवातीला वळून. moved_keys(keys, before, after) ज्यांचा मालक बदलला त्या keys मोजते. design/demo.py मधील balance() 10,000 वर्ग keys दोन्ही पद्धतींनी 4 वरून 5 servers वर हलवते आणि प्रत्येक server वरच्या keys छापते.

🧪 करून पाहा

python3 design/demo.py balance
python3 - <<'EOF'
import sys; sys.path.insert(0, "design"); from blocks import HashRing, moved_keys
keys = [f"class-{i}" for i in range(10_000)]
for v in (1, 10, 100):
    ring = HashRing(["s0", "s1", "s2", "s3"], vnodes=v)
    load = {}
    for k in keys: load[ring.owner(k)] = load.get(ring.owner(k), 0) + 1
    print(f"vnodes {v:>3} → busiest server {max(load.values()):>5,} keys · quietest {min(load.values()):>5,}")
four, five = HashRing(["s0", "s1", "s2", "s3"]), HashRing(["s0", "s1", "s2", "s3", "s4"])
eight = HashRing([f"s{i}" for i in range(8)])
print(f"ring 4 → 8 servers: {moved_keys(keys, four.owner, eight.owner):,} keys move")
print(f"ring 5 → 4 (s4 dies): {moved_keys(keys, five.owner, four.owner):,} keys move")
h = lambda k: int(k.split('-')[1]) * 2654435761 % 2**32
for n in (4, 10, 100):
    print(f"hash % N, {n:>3} → {n + 1:>3} servers: {moved_keys(keys, lambda k: h(k) % n, lambda k: h(k) % (n + 1)) / 100:.0f}% move")
EOF

✅ तपासा — तुम्हाला काय दिसायला हवे

balance छापते:

── 10,000 classes on 4 cache servers; a 5th server joins
   hash % N         → 8,001 keys move (80%)
   consistent hash  → 1,975 keys move (20%) — about 1/5, only the new server's share
   keys per server with 100 virtual nodes each: {'s0': 1883, 's1': 2016, 's2': 1989, 's3': 2137, 's4': 1975}
   load balancers: L4 (TCP, fast) vs L7 (HTTP paths, headers) · health checks · round robin / least connections

तुमचा snippet छापतो:

vnodes   1 → busiest server 4,576 keys · quietest   184
vnodes  10 → busiest server 3,358 keys · quietest 2,012
vnodes 100 → busiest server 2,585 keys · quietest 2,282
ring 4 → 8 servers: 4,967 keys move
ring 5 → 4 (s4 dies): 1,975 keys move
hash % N,   4 →   5 servers: 80% move
hash % N,  10 →  11 servers: 91% move
hash % N, 100 → 101 servers: 99% move

🏁 तुम्ही आत्ताच काय सिद्ध केले

% N सह, ताफा जितका मोठा तितके वाईट: 10 → 11 मध्ये 91% हलतात, 100 → 101 मध्ये 99%. Ring फक्त नव्या server चा वाटा हलवते (1,975 = जोडल्यानंतर s4 च्या मालकीच्या नेमक्या तितक्याच keys), आणि 4 → 8 दुप्पट केल्यावर सुमारे अर्ध्या हलतात — नवे servers आता ज्यांचे मालक आहेत त्या अर्ध्या. s4 गेला की त्याच 1,975 keys परत हलतात, बाकी एकही नाही. प्रत्येक server ला एक बिंदू असताना सर्वात व्यस्त server कडे 4,576 keys आणि सर्वात शांत server कडे 184; 100 virtual nodes सह, 2,585 vs 2,282.

⚠️ नेहमीच्या चुका

🏭 प्रत्यक्ष वापरात

On a real account — health check आणि least-outstanding algorithm असलेला ALB target group, Terraform मध्ये:

resource "aws_lb_target_group" "api" {
  name                          = "notice-api"
  port                          = 8080
  protocol                      = "HTTP"
  vpc_id                        = var.vpc_id
  load_balancing_algorithm_type = "least_outstanding_requests"
  deregistration_delay          = 30            # let in-flight requests finish
  health_check {
    path                = "/health"
    interval            = 10
    healthy_threshold   = 2
    unhealthy_threshold = 3
    matcher             = "200"
  }
}

तुम्ही स्वतः hash ring क्वचितच लिहिता: Redis Cluster keys ना 16,384 hash slots मध्ये विभागतो आणि तुम्ही shard जोडला की slots हलवतो; DynamoDB आणि Cassandra तुमच्यासाठी key च्या hash नुसार partition करतात. तुमचे काम म्हणजे पसरणारी key निवडणे ("आजची तारीख" नाही).

अधिक खोलात: Scaling शाळा, धडे 05–07 (stateless servers, load balancers, Auto Scaling) आणि 11 (partitioning आणि hot keys).

🏭 प्रत्यक्ष वापरात हे का महत्त्वाचे: capacity वाढवायची योजना करताना विचारा "काय हलते?" उत्तर "जवळजवळ सगळे" असेल, तर बदल शांत तासासाठी ठरवा — किंवा आधी placement चा नियम बदला.

⏭️ पुढे

काही काम पालक थांबलेला असताना करण्यासाठी खूप हळू असते — जसे 5 million SMS पाठवणे. नियोजन कार्यालय त्याला एक queue देते.

git checkout lesson-07-queues

⚖️ Lesson 06 — Load balancing + consistent hashing: spread the visitors and the keys

📍 You are here: Lesson 06 of 18 · Previous: lesson-05-caching · Next: lesson-07-queues


📦 What's in this branch

Lessons 01–05, plus two kinds of spreading. Load balancers spread requests over identical API servers: L4 vs L7, algorithms, health checks. Consistent hashing spreads keys over cache (or data) servers so that adding or losing one server moves only its share. HashRing and moved_keys() in design/blocks.py count the keys that move; balance() in design/demo.py compares hash % N with the ring.

🧒 Explain like I'm 5

The gate helper. On a busy morning 12 clerks sit at 12 counters. A helper at the gate sends each parent to a counter: the next one in turn, or the one with the shortest queue. Every few minutes she checks each counter — "are you open?" — and stops sending parents to a clerk who went home. That helper is a load balancer.

The pigeon-holes. Now a different job. The office keeps copies of class notices in 4 cupboards. Which cupboard holds class 3A? Dipika's first rule: "class number divided by 4 — the remainder is the cupboard." It works. Then a 5th cupboard arrives. Now the rule is "divided by 5", and almost every class lands in a different cupboard. The clerks must move 8 of every 10 folders. All the copies are "lost" at once.

Dipika's better rule: paint a big clock face ⏰ on the wall. Each cupboard gets a few marks around the clock. Each class also gets a spot on the clock, and it lives in the next cupboard mark going clockwise. When a 5th cupboard arrives, it adds its marks — and takes only the classes just before its marks. About 1 in 5 folders move. The rest stay where they are.

🗺️ Diagram

flowchart LR
    subgraph lb["⚖️ load balancer (requests)"]
      c["👪 requests"] --> alb["L7 ALB<br/>round robin / least conn<br/>health checks"]
      alb --> a1["API 1"]
      alb --> a2["API 2"]
      alb --> a3["API …12"]
    end
    subgraph ring["⏰ consistent hashing (keys)"]
      k["10,000 class keys"] --> m["hash % N, 4 → 5<br/>8,001 move (80%)"]
      k --> r["hash ring, 4 → 5<br/>1,975 move (20%)<br/>100 virtual nodes each"]
    end
    a1 --> r

🗺️ Drawn version + a lab: https://school-edh.pages.dev/system-design/lesson-diagrams.html#l06

❓ What

🤔 Why

Because lesson 02 needs 12 API servers and a cache that may grow from 4 nodes to 5. The load balancer makes 12 servers look like one and hides a failed one. Consistent hashing makes the cache grow without a cliff: with % N, adding a node on a busy morning empties 80% of the cache and sends those reads to the database at the worst moment.

🔧 How (in this repo)

HashRing(servers, vnodes=100) in design/blocks.py puts vnodes points per server on the ring (MD5 of "s0#17" …), kept sorted with bisect. owner(key) finds the first point after the key's hash, wrapping around. moved_keys(keys, before, after) counts keys whose owner changed. balance() in design/demo.py moves 10,000 class keys from 4 to 5 servers both ways and prints the keys per server.

🧪 Try it

python3 design/demo.py balance
python3 - <<'EOF'
import sys; sys.path.insert(0, "design"); from blocks import HashRing, moved_keys
keys = [f"class-{i}" for i in range(10_000)]
for v in (1, 10, 100):
    ring = HashRing(["s0", "s1", "s2", "s3"], vnodes=v)
    load = {}
    for k in keys: load[ring.owner(k)] = load.get(ring.owner(k), 0) + 1
    print(f"vnodes {v:>3} → busiest server {max(load.values()):>5,} keys · quietest {min(load.values()):>5,}")
four, five = HashRing(["s0", "s1", "s2", "s3"]), HashRing(["s0", "s1", "s2", "s3", "s4"])
eight = HashRing([f"s{i}" for i in range(8)])
print(f"ring 4 → 8 servers: {moved_keys(keys, four.owner, eight.owner):,} keys move")
print(f"ring 5 → 4 (s4 dies): {moved_keys(keys, five.owner, four.owner):,} keys move")
h = lambda k: int(k.split('-')[1]) * 2654435761 % 2**32
for n in (4, 10, 100):
    print(f"hash % N, {n:>3} → {n + 1:>3} servers: {moved_keys(keys, lambda k: h(k) % n, lambda k: h(k) % (n + 1)) / 100:.0f}% move")
EOF

✅ Verify — what you should see

balance prints:

── 10,000 classes on 4 cache servers; a 5th server joins
   hash % N         → 8,001 keys move (80%)
   consistent hash  → 1,975 keys move (20%) — about 1/5, only the new server's share
   keys per server with 100 virtual nodes each: {'s0': 1883, 's1': 2016, 's2': 1989, 's3': 2137, 's4': 1975}
   load balancers: L4 (TCP, fast) vs L7 (HTTP paths, headers) · health checks · round robin / least connections

Your snippet prints:

vnodes   1 → busiest server 4,576 keys · quietest   184
vnodes  10 → busiest server 3,358 keys · quietest 2,012
vnodes 100 → busiest server 2,585 keys · quietest 2,282
ring 4 → 8 servers: 4,967 keys move
ring 5 → 4 (s4 dies): 1,975 keys move
hash % N,   4 →   5 servers: 80% move
hash % N,  10 →  11 servers: 91% move
hash % N, 100 → 101 servers: 99% move

🏁 What you just proved

With % N, the bigger the fleet, the worse it gets: 10 → 11 moves 91%, 100 → 101 moves 99%. The ring moves only the new server's share (1,975 = exactly the keys s4 owns after joining), and doubling 4 → 8 moves about half — the half the new servers now own. When s4 dies, the same 1,975 keys move back, and no others. With one point per server the busiest one holds 4,576 keys and the quietest 184; with 100 virtual nodes, 2,585 vs 2,282.

⚠️ Common mistakes

🏭 In production

On a real account — an ALB target group with a health check and the least-outstanding algorithm, in Terraform:

resource "aws_lb_target_group" "api" {
  name                          = "notice-api"
  port                          = 8080
  protocol                      = "HTTP"
  vpc_id                        = var.vpc_id
  load_balancing_algorithm_type = "least_outstanding_requests"
  deregistration_delay          = 30            # let in-flight requests finish
  health_check {
    path                = "/health"
    interval            = 10
    healthy_threshold   = 2
    unhealthy_threshold = 3
    matcher             = "200"
  }
}

You rarely write a hash ring yourself: Redis Cluster splits keys into 16,384 hash slots and moves slots when you add a shard; DynamoDB and Cassandra partition by a hash of the key for you. Your job is to pick a key that spreads (not "today's date").

Go deeper: the Scaling school, lessons 05–07 (stateless servers, load balancers, Auto Scaling) and 11 (partitioning and hot keys).

🏭 Why this matters in production: when you plan to add capacity, ask "what moves?" If the answer is "almost everything", plan the change for a quiet hour — or change the placement rule first.

⏭️ Next

Some work is too slow to do while the parent waits — like sending 5 million SMS. The planning office gives it a queue.

git checkout lesson-07-queues
← PreviouscachingNext →queues

This page is the lesson's README from the lesson-06-load-balancing branch, shown here so the whole School stays on one site. Code files open on GitHub at the same branch.