Load Balancing and Traffic Management Questions
Distributing requests across capacity: load-balancing algorithms (round-robin, least-connections, consistent hashing), L4 versus L7 balancing, health checks, and traffic shaping. Covers sticky sessions, canary and blue-green routing, rate limiting, and graceful draining. The traffic-distribution layer that keeps a scaled system balanced and available.
Compare consistent hashing and rendezvous (highest-random-weight) hashing. How does each handle node addition and removal, what is the cost of rebalancing, and which would you use for cache routing versus partitioned storage with heterogeneous node weights?
Sample Answer
Direct Answer
Consistent hashing (CH) places nodes and keys on a hash ring and routes each key to the next node clockwise; rendezvous hashing (also called highest-random-weight, HRW) scores every node for a given key with an independent hash and routes to the node with the highest score. Both move only a small, bounded fraction of keys when the node set changes, but they differ in what they need to operate: CH needs a maintained ring data structure (and virtual nodes to balance load), while HRW needs nothing shared at all, it recomputes from scratch on every lookup. That operational difference, not raw performance, usually decides which one you pick.
How Each Handles Node Changes
Consistent hashing. Nodes (or their virtual tokens) occupy positions on a ring. A key is owned by the first node clockwise of its hash position. Adding a node inserts a new position and reassigns only the keys that fall in the arc between the new node and its counterclockwise neighbor. Removing a node deletes its position, and in a plain single-token ring, every key that node owned is reassigned entirely to its clockwise successor. That successor absorbs the whole load, it is not spread across the rest of the cluster. This is why production consistent-hash rings almost always use virtual nodes: giving each physical node many tokens scattered around the ring means a removal's load lands on many different successors instead of one.
Rendezvous hashing. For a key, each node's score is an independent draw from a hash function, and the node with the maximum score wins. Adding a node adds one more independent score into the comparison; only keys where the new score happens to be the maximum move, and they move only to the new node. Removing a node drops one score out of the comparison; for every key that node used to win, the new winner is whichever of the remaining nodes had the second-highest score, which, because all scores are independently drawn from the same distribution, is uniformly distributed across the survivors. Rendezvous hashing gets even redistribution on removal for free, without virtual nodes.
Cost of Rebalancing
| Operation | Consistent hashing (single-token) | Consistent hashing (with V virtual nodes/node) | Rendezvous hashing |
|---|---|---|---|
| Node join | Only the new node's arc moves | Only the new node's arc(s) move | Only keys where the new score wins move |
| Node leave | All of the removed node's keys land on one successor | Removed node's keys spread across roughly V nearby successors | Removed node's keys spread uniformly across all survivors |
| Lookup cost | O(log T) with T ring tokens (binary search) | O(log T), T grows with V | O(N) hash evaluations per key, N = node count |
| Metadata to maintain | Sorted ring of N positions | Sorted ring of N*V positions | None beyond the current node list |
Worked Example
Take a cluster of N=10 nodes holding keys uniformly on the ring, so each node owns an expected 101 of the keyspace.
Join (both schemes). Adding one node makes it 1 of 11 exchangeable ring owners (CH) or 1 of 11 exchangeable score-holders (HRW). By symmetry, the new node's expected share is:
N+11=111≈9.1%and every other node's share is unaffected in expectation, for both schemes.
Leave, plain single-token consistent hashing. The removed node held N1=101=10% of keys. All of it moves to one clockwise successor, whose share becomes:
101+101=102=20%That successor now carries double its prior load, a hotspot, while the other 8 surviving nodes are untouched.
Leave, rendezvous hashing (or CH with enough virtual nodes to approximate it). The removed node's 10% share spreads uniformly across the 9 survivors. Each survivor's new share:
101+10×91=909+1=91≈11.1%which is exactly the uniform N−11 share you would expect if the whole keyspace were freshly and evenly split across 9 nodes. This algebra generalizes: N1+N(N−1)1=N−11, confirming that rendezvous hashing (and a well-provisioned virtual-node ring) reach the ideal uniform post-removal distribution automatically, while a plain single-token ring does not.
Which to Use
For cache routing (many roughly similar-capacity nodes, request-level lookups, simplicity matters, weights may differ per node's memory size), rendezvous hashing is usually the simpler choice: no ring to maintain, exact per-node weighting by scaling scores, and even load redistribution on failure without tuning a virtual-node count. The O(N) lookup cost is rarely the bottleneck unless N reaches the high hundreds or thousands of cache nodes.
For partitioned storage with a large number of partitions (sharded databases, distributed key-value stores), consistent hashing with virtual nodes is usually preferred: ownership maps to contiguous ranges, which supports range scans and simple "move this token's range" rebalancing tooling, and lookup stays O(log T) even as the ring grows large. Heterogeneous weights are handled by giving higher-capacity nodes proportionally more virtual nodes, though that is a discrete approximation rather than an exact ratio.
Trade-offs and Pitfalls
- The single-token consistent-hashing hotspot on removal (shown above) is a frequent interview trap: candidates who describe CH without mentioning virtual nodes are implicitly describing a scheme with an uneven-failover problem.
- Rendezvous hashing's O(N) per-lookup cost is a real scaling limit for very large node counts; production systems mitigate it with a hierarchical or bounded-candidate variant rather than switching to CH purely for that reason.
- Weighting is exact and continuous with rendezvous hashing (scale each node's score by its weight) but only approximate with CH virtual nodes (weight is quantized by how many tokens you assign).
- Neither scheme is free of the "every key touches every node during full re-sharding" cost if you ever need to change the hash function itself (e.g., migrating hash algorithms); that requires a coordinated dual-write or shadow-ring migration regardless of which scheme you use.
Explain active versus passive health checks used by load balancers and service discovery. For each, describe typical probe frequency, example probes for HTTP, gRPC, and TCP, and how each affects failover decisions. What's a good strategy for combining them in production to reduce false positives?
Sample Answer
Direct answer
Active checks are probes the load balancer or service discovery system initiates on a schedule; passive checks are inferences drawn from real request outcomes, such as a 5xx response, a timeout, or a connection reset, as traffic naturally flows. Active checks can catch a problem before a real user hits it, but can false-positive on transient blips; passive checks add no extra probe traffic and reflect real user impact directly, but by definition only detect a failure after at least one real request has already suffered it. Production systems combine both.
Structured elaboration
| Active checks | Passive checks | |
|---|---|---|
| Typical frequency | 1 to 10 seconds for load-balancer probes, 5 to 30 seconds for service-discovery health endpoints | Continuous, evaluated over a rolling window of real requests |
| HTTP example | GET a health endpoint, expect a 200 with an optional body pattern | Observed 5xx responses or request timeouts on real traffic |
| gRPC example | Call the standard grpc.health.v1.Health service, expect SERVING | Observed non-OK status codes or stream resets on real calls |
| TCP example | Connect and optionally read an application banner | Observed connection resets or handshake failures on real connections |
| Detects | Problems visible to a synthetic probe, proactively | Problems that actually affect real traffic, reactively |
| Failover trigger | N consecutive probe failures | Error rate or failure count over a rolling window of real requests |
Worked example
A production strategy that combines both without over-reacting to a single blip: mark an instance "suspect" after 2 consecutive active probe failures, but do not remove it from rotation yet. Corroborate with a passive signal over a rolling window of the last 20 real requests: if 50% or more of those requests failed, that is a threshold of 20×0.5=10 failed requests out of 20, remove the instance from rotation. Requiring both signals means a single flaky probe or a single unlucky real request cannot remove a healthy instance on its own, while a genuinely failing instance still gets caught quickly because the two signals reinforce each other.
Trade-offs & pitfalls
- Active-only checks can pass while real traffic fails: a backend can answer a generic TCP or HTTP probe correctly while returning errors to requests carrying real auth headers or payloads the probe never sends.
- Passive-only checks mean, by definition, that some real users experience the failure before it is detected; that is not acceptable on its own for canary or rollout gating, where the goal is to catch a bad release before it reaches most users.
- Avoid routing probe traffic through the same accounting a load-balancing algorithm like least connections uses for real capacity; counting probes as connections can skew the balance away from actual user traffic.
Design a monitoring and alerting setup for a load balancer tier so you catch unhealthy backends before users see errors. Which metrics would you include and why, what alert thresholds would you propose, and how would you differentiate severity to avoid alert fatigue?
Sample Answer
Direct answer
Catching unhealthy backends before users notice means monitoring the LB tier along four angles: traffic, latency, and errors as seen from the LB itself (the RED metrics), backend health signals independent of user-facing symptoms, saturation and capacity of the LB tier, and TLS/connection-layer health, then tying alert severity to how much of the user-facing error budget (the amount of failure your reliability target still allows before you've broken your promise to users) is actually being consumed rather than to a single static threshold, so a brief blip and a sustained outage don't page the same way.
The metric set
| Category | Metric | Why it matters |
|---|---|---|
| Traffic | Requests per second, overall and per backend pool | Baseline for interpreting every other metric; a rate change can explain an error-rate change |
| Latency | p50 / p90 / p95 / p99, split into LB-processing time and end-to-end time | Percentiles catch tail degradation that an average hides; splitting LB time from backend time tells you where the slowdown is |
| Errors | 4xx vs 5xx rate, and whether 5xx originated at the LB (no healthy backend) or was passed through from a backend | A 5xx-at-the-LB spike usually means capacity loss, not an application bug |
| Backend health | Per-backend health-check pass/fail rate and consecutive-failure counts | The direct signal for whether a specific node is about to be pulled |
| Connections | Active connections and connection churn (opens/closes per second) | A churn spike often precedes a latency or error spike |
| Capacity | LB-tier CPU, memory, socket/file-descriptor usage | The LB tier itself can become the bottleneck, not just the backends |
| TLS | Handshake failure rate, certificate expiry countdown | Silent failure mode: expired certs or handshake issues look like "backend is fine but nobody can connect" |
| Traffic shaping | Rate-limit rejections, circuit-breaker open/close events | Leading indicators: a circuit breaker tripping on a backend is an early warning before that backend fails health checks outright |
Alert severity: tie it to error-budget burn, not just a raw threshold
A raw threshold ("page if 5xx rate > 1%") treats a 10-minute blip the same as a slow multi-hour decline, a common cause of alert fatigue in either direction: it either pages on noise or is set so loose it misses real problems. The more robust approach expresses the threshold as a multiple of your error-budget burn rate: if you're consuming your monthly error budget faster than the rate at which "consume it evenly over 30 days" would, page proportionally to how much faster.
Worked example
Suppose the LB tier's SLO (service-level objective, the reliability target you've committed to, here 99.9% of requests succeeding) is 99.9% success rate over a 30-day window. The allowed error budget is 0.1% of requests, equivalently, 0.1% of the 30-day window can be fully down and still meet the SLO:
- 30 days in minutes: 30×24×60=43,200 minutes.
- Budget: 43,200×0.001=43.2 minutes of full-downtime-equivalent for the entire month.
Now suppose the current 5xx rate over the last hour is holding at 2%, against an objective of 0.1%. The burn rate is:
burn rate=0.1%2%=20×At 20 times the sustainable rate, the entire month's budget would be exhausted in:
2030 days=1.5 daysA burn rate that would exhaust the monthly budget in 1.5 days is unambiguously page-worthy; a burn rate of, say, 1.2 times sustainable is not, it just needs tracking. That is the difference between a severity-1 page and a low-priority ticket, expressed in the same unit (budget consumption), not two arbitrarily different percentage cutoffs.
Avoiding alert fatigue
- Tier severity by projected budget-exhaustion time (minutes-to-exhaust pages immediately, weeks-to-exhaust becomes a ticket), rather than by a single static percentage.
- Deduplicate and correlate: a health-check-flapping alert and a 5xx-rate alert firing from the same underlying cause should collapse into one incident, not page twice.
- Route by audience: on-call needs the actionable, current-state view (which backend, since when); a weekly review needs the trend (is health-check-induced churn increasing over time).
- Require every page-level alert to link directly to a runbook step; an alert with no clear next action trains people to ignore it.
Trade-offs and pitfalls
Over-indexing on system-level metrics (CPU, connections) without a user-facing correlate risks paging on things users never notice; conversely, only alerting on aggregate error rate can hide a problem affecting one backend pool or one tenant while the aggregate stays healthy. Circuit-breaker and rate-limit rejection counts are easy to skip because they're not "errors" in the traditional sense, but they're often the earliest signal, by the time health checks start failing outright, the circuit breaker has usually already tripped.
Implement a consistent hashing ring that supports weighted, heterogeneous-capacity backends: add_node(node_id, weight), remove_node(node_id), and get_node(key). Explain how you map weight to a number of virtual nodes without creating an excessive number of them for very large weights, and show a small example demonstrating minimal key movement when a node is added or removed.
Sample Answer
Direct answer
Map weight to virtual-node count linearly (replicas proportional to weight), because that is what preserves the property that a node's expected key share equals its share of total weight, and cap the replica count per node so one extreme outlier weight cannot blow up ring size or per-lookup cost. Adding or removing a node then only remaps the keys that specifically belonged to that node's virtual replicas, everyone else's keys stay put, which is the whole point of consistent hashing over a plain hash-mod-N scheme.
Approach
The core mechanism is the same ring as unweighted consistent hashing (sorted virtual-node positions, binary search for ownership), with two additions: add_node takes a weight, and the number of virtual replicas it creates is proportional to that weight rather than a fixed constant. This keeps a node's fraction of total ring positions equal to its fraction of total weight, so its expected fraction of keys matches too (this is verified numerically below).
Without a cap, a single node with weight 100000 relative to peers at weight 1 would ask for 100000 times the base replica count, which is both a memory problem (ring size) and a lookup problem (O(log n) grows with total replicas). The cap trades strict weight-proportionality at the extreme end for a bounded ring: past the cap, an outlier-weighted node still gets more replicas than its peers, just not linearly more, which is an acceptable trade since a weight ratio that extreme usually means the two backends are different enough in kind that the whole model, not just the replica formula, should be reconsidered.
Code (Python)
import bisect
import hashlib
class ConsistentHashRing:
def __init__(self, base_replicas=100, max_replicas=4000):
self.ring = [] # sorted hash positions
self.nodes = {} # position -> node_id
self.node_positions = {} # node_id -> list of positions (for O(replicas) removal)
self.base_replicas = base_replicas
self.max_replicas = max_replicas # caps ring growth for very large weights
def _hash(self, key: str) -> int:
return int(hashlib.sha256(key.encode("utf-8")).hexdigest(), 16)
def _replica_count(self, weight: float) -> int:
# Linear in weight: a node's replica share stays proportional to its
# weight share, which is what makes expected key share == weight share.
# max_replicas bounds ring memory (O(replicas)) and lookup cost
# (O(log total_replicas)) against one extreme outlier weight.
raw = int(round(self.base_replicas * weight))
return max(1, min(raw, self.max_replicas))
def add_node(self, node_id: str, weight: float = 1.0):
replicas = self._replica_count(weight)
positions = []
for i in range(replicas):
h = self._hash(f"{node_id}#{i}")
if h in self.nodes:
continue
bisect.insort(self.ring, h)
self.nodes[h] = node_id
positions.append(h)
self.node_positions[node_id] = positions
def remove_node(self, node_id: str):
for h in self.node_positions.pop(node_id, []):
idx = bisect.bisect_left(self.ring, h)
if idx < len(self.ring) and self.ring[idx] == h:
self.ring.pop(idx)
self.nodes.pop(h, None)
def get_node(self, key: str):
if not self.ring:
return None
h = self._hash(key)
idx = bisect.bisect_right(self.ring, h)
if idx == len(self.ring):
idx = 0
return self.nodes[self.ring[idx]]
Key points
- Weight-to-replica mapping is linear, then capped:
replicas = clamp(round(base_replicas * weight), 1, max_replicas). Linear preserves proportional key share; the cap bounds ring size. node_positionstracks each node's own replica hashes soremove_nodeonly has to touch that node's O(replicas) entries, not scan the whole ring.- Minimal movement is structural, not incidental: because virtual nodes for different physical nodes are interleaved around the ring, adding one new node only steals the ring segments immediately preceding its own virtual points, wherever those land, from whichever nodes currently own them; it cannot affect a segment it doesn't touch.
Complexity
add_node: O(replicas * log n) where replicas is capped atmax_replicasand n is total virtual nodes on the ring.remove_node: O(replicas * log n) for that node's own replicas only.get_node: O(log n) for the binary search.- Space: O(n) total virtual nodes across all live physical nodes, bounded by (number of nodes) times
max_replicas. - Caveat: the O(log n) parts above cover only the position search.
bisect.insortandlist.pop(idx)mutate a plain Python list, so the actual insertion/removal is O(n) per call (array shift), not O(log n); the true cost ofadd_node/remove_nodeis O(replicas * n) in the worst case. A balanced tree or skip list would be needed to make the mutation itself sub-linear at large node counts.
Worked example (weight-to-replica mapping and minimal movement)
Weight-to-replica mapping with base_replicas=100, max_replicas=4000, showing the cap engage at large weight:
weight= 1 -> replicas=100
weight= 4 -> replicas=400
weight= 16 -> replicas=1600
weight= 100 -> replicas=4000
weight= 100000 -> replicas=4000
Weight 100 already saturates the cap at this configuration (100 x 100 = 10000, clamped to 4000); weight 100000 hits the same cap, confirming the ring size stays bounded regardless of how extreme a weight is declared.
Distribution check with three nodes of weight 1, 4, and 1 (total weight 6), over 20000 fixed keys (key-0 through key-19999):
weighted distribution (A:1, B:4, C:1), 20000 keys: {'A': 3483, 'B': 13097, 'C': 3420}
expected B share ~= 4/6 = 0.6667, measured = 0.6549
B's weight share is 4/6 = 0.6667; its measured key share is 0.6549, within normal sampling variance for a ring at this replica density.
Minimal-movement check: adding a fourth node D with weight 2 (new total weight 8):
expected moved fraction=total weight afterweightD=82=0.25keys moved after adding D(weight=2): 4583/20000 = 0.2291
of those, moved specifically to D: 4583/4583 = 1.0000 (should be ~1.0)
22.91% of keys moved, close to the 25% expectation, and (critically) every single key that moved, moved specifically to D, none of A/B/C's other keys were disturbed by each other. Removing B afterward confirms the same locality in the other direction:
after removing B: 9957 keys needed reassignment, 0 non-B keys moved
Every key that needed reassignment had previously belonged to B; zero keys that belonged to A, C, or D were touched by B's removal, this is the "minimal key movement" property the question asks to demonstrate, shown numerically rather than asserted.
Edge cases
- Weight rounds down to 0 replicas: guarded with
max(1, ...)so every added node gets at least one ring position, even a very small declared weight. - Extremely large weight: bounded by
max_replicas, verified above. - Removing a node not present: no-op, since
node_positions.pop(node_id, [])returns an empty list. - Empty ring:
get_nodereturnsNone.
Trade-offs and pitfalls
- The cap breaks strict proportionality at the high end. Two nodes both past the cap (say weight 500 and weight 100000) get the same replica count and therefore roughly the same key share, even though their declared weights differ by 200x; if that distinction matters operationally, the cap needs to be raised or the weight scale needs to be redefined so realistic weights don't approach it.
base_replicasandmax_replicasare a joint tuning knob, not independent ones. Raisingbase_replicasto smooth distribution variance also lowers the effective weight ceiling before the cap engages (sincemax_replicas / base_replicasis the highest weight that still gets proportional treatment); changing one without reconsidering the other can silently shrink your usable weight range.- Weight changes at runtime are a remove-then-add, not an update-in-place, in this design; recomputing a node's replica count means discarding its old positions and generating new ones, which moves keys even for a node that didn't fail, this is worth calling out explicitly if the interviewer asks about live capacity rebalancing.
- Ring divergence across processes is still the same risk as unweighted consistent hashing: every client needs the same node set, weights,
base_replicas,max_replicas, and hash function to agree on ownership.
How does weighted round robin differ from applying weights to least-connections? Sketch how you would distribute requests proportional to backend weight under each approach, and describe when you would adjust weights dynamically, for example during autoscaling or when an instance is degraded.
Sample Answer
Direct answer
Both are ways to make an algorithm capacity-aware, but they optimize different things. Weighted round robin (WRR) distributes a fixed pattern of requests proportional to weight, without looking at current load: it is a scheduling problem, solved up front. Weighted least connections picks, per request, whichever server has the lowest connections-to-weight ratio: it is a live load-balancing decision that reacts to what's actually happening on each server right now. WRR is cheaper and predictable; weighted least connections is more adaptive when request duration varies.
Weighted round robin mechanics
A naive WRR just repeats each server in the rotation list a number of times equal to its weight (weight 5, 1, 1 becomes the list [A,A,A,A,A,B,C]), which is simple but bursty, all of A's requests land in a clump. The standard fix is smooth weighted round robin, used by nginx and others: each server keeps a running counter, and at every request:
then the server with the highest counter is selected, and only that server's counter is reduced by the total weight:
ci∗(t+1)=ci(t+1)−W,W=j∑wjWorked trace: weights A=5, B=1, C=1 (W=7)
| Step | Counters after add (A, B, C) | Selected | Counters after subtract |
|---|---|---|---|
| 1 | 5, 1, 1 | A | -2, 1, 1 |
| 2 | 3, 2, 2 | A | -4, 2, 2 |
| 3 | 1, 3, 3 (tie B/C, broken alphabetically) | B | 1, -4, 3 |
| 4 | 6, -3, 4 | A | -1, -3, 4 |
| 5 | 4, -2, 5 | C | 4, -2, -2 |
| 6 | 9, -1, -1 | A | 2, -1, -1 |
| 7 | 7, 0, 0 | A | 0, 0, 0 |
Sequence: A, A, B, A, C, A, A. Final counts: A=5, B=1, C=1, exactly matching the declared weights, and after 7 steps every counter returns to 0, so the pattern repeats cleanly. Compare this to the naive list [A,A,A,A,A,B,C]: smooth WRR interleaves B and C between A's turns instead of clumping them at the end.
Weighted least connections mechanics
Instead of a precomputed pattern, each request goes to whichever server minimizes:
scorei=wiciwhere ci is current active connections. This needs live state (the balancer must track connection counts accurately) but it self-corrects when request durations vary: if one of A's requests is unusually slow and its connection count stays elevated, weighted least connections routes around it immediately, while WRR would keep sending A its scheduled 5-out-of-7 share regardless.
When to adjust weights dynamically
- Autoscaling: when an instance joins or leaves the pool, recompute weights from the group's new capacity (commonly proportional to vCPU count or a load-tested throughput figure), otherwise new instances sit idle while old ones stay pinned to stale weights.
- Degraded instance: on failed or borderline health checks, reduce the instance's weight toward zero (a soft drain) rather than hard-removing it, so in-flight requests finish before it stops receiving new ones.
- Heterogeneous hardware: base weights on measured throughput under load, not just instance-type labels, since two "same size" instances can have different real capacity due to noisy neighbors or different hardware generations.
- Oscillation risk: reacting to every latency blip by changing weights can cause a feedback loop (a server's weight drops, it gets less traffic, its metrics improve, weight goes back up, traffic returns, metrics degrade again). Apply a decay or minimum-dwell-time before a weight change takes effect.
Trade-offs and pitfalls
WRR is stateless and cheap to run at very high request rates because it never inspects live connection counts, but it is blind to reality: if a server silently starts responding slowly, WRR keeps sending it its full scheduled share. Weighted least connections fixes that but needs accurate, low-latency visibility into connection counts across the fleet, which is harder in a distributed balancer with multiple LB instances that don't share state. A common middle ground is weighted least connections with a smoothing factor so weight changes and connection counts don't cause the oscillation described above.
Unlock Full Question Bank
Get access to all Load Balancing and Traffic Management interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.