Skip to content

✨ Generated Additions

Distributed Lock Service Design

A distributed lock service provides mutual exclusion across multiple nodes. The core challenges are handling network partitions, node failures, and preventing two clients from holding the same lock simultaneously (split-brain).

Core Building Blocks

  • Strongly consistent store: Use a consensus-based coordination service such as ZooKeeper, etcd, or Chubby. These systems provide linearizable writes and sequential consistency, which are required for safe lock acquisition. Avoid using plain Redis without careful fencing because a single Redis instance is a single point of failure, and Redlock-style multi-node algorithms can still lose safety under GC pauses or network delays.
  • Lease-based locks: Each lock is granted with a lease (TTL). The client must renew the lease before it expires. If the client crashes or is partitioned, the lock is automatically released after the TTL, allowing other clients to make progress.
  • Fencing tokens: When a lock is acquired, the lock service returns a monotonic token (e.g., a version number). The client passes this token with every operation. Storage or other services reject operations with stale tokens, preventing a delayed client that still believes it holds the lock from causing corruption.
  • Lock metadata: Store lock holder ID, lease expiration, and fencing token in a persistent, replicated path (e.g., a ZooKeeper znode or etcd key). Clients use atomic compare-and-swap or create-with-uniqueness to acquire locks.

How This Answers the Question

When designing a distributed lock service, explicitly choose a consensus-based backend for leadership and linearizability. Describe the acquisition / renewal / release protocol, the role of leases to tolerate crashes, and fencing tokens to guard against split-brain. Mention trade-offs: a lock service built on ZooKeeper/etcd favors safety over raw throughput; a Redis-based approach is faster but must include fencing and careful failure analysis. This explains how the lock service coordinates mutually exclusive access across distributed clients safely and with high availability.


Scaling the Load Balancer Itself: Distributed LB Architecture

A single load balancer is a bottleneck and single point of failure. To design a distributed load balancer, you must horizontally scale the load-balancing layer itself and coordinate multiple LB instances.

Core ideas

  • Multiple LB nodes behind a virtual IP (VIP) or anycast address. DNS resolves the LB name to multiple IPs (or anycast routes clients to the nearest LB node).
  • Stateless LB nodes so any node can handle any request. Persistent state (e.g., session stickiness tables) is stored in a fast distributed store like Redis or etcd, not locally.
  • Consistent hashing to distribute traffic across backend servers while minimizing rehashing when backend servers are added or removed. This also helps LBs agree on which backend should receive a given request without sharing all state.
  • Coordination via a consensus service (ZooKeeper, etcd) for cluster membership, leader election (if needed), and health status of backend servers. Each LB node subscribes to updates and maintains a consistent view.
  • Data plane vs control plane separation: the data plane (fast packet forwarding, often in kernel or eBPF) is distributed across LB nodes; the control plane (configuration, health checks, routing policies) is managed centrally and pushed to all nodes.
  • Health checks at scale: each LB node periodically probes backends and shares results with peers (or via control plane). If a node dies, others take over its traffic (via ECMP, anycast failover, or DNS repointing).
  • Avoiding split-brain: for active-active LB pairs, use consistent hashing and shared state to ensure only one node acts on a given request; for active-passive, use a leader-elect algorithm (e.g., via etcd lease) so only the leader owns the VIP.

How this answers the question

When asked to design a distributed load balancer, this section explains how to make the LB itself scalable and highly available: stateless nodes, consistent hashing, shared coordination, and control/data plane split. It addresses the key challenges: avoiding a single point of failure, distributing load evenly, maintaining session affinity when needed, and performing health checks across a large backend fleet. Use this to describe a concrete architecture (e.g., anycast + ECMP to reach multiple L4 LB nodes, which then route to L7 proxies using consistent hashing).


Full Walkthrough — Design a Chat App like WhatsApp

Requirements and constraints

Start by clarifying the scope: one-on-one messages, group messages, media attachments, online presence, read receipts, multi-device sync, and offline delivery. The main non-functional requirements are low latency for message delivery, reliable and ordered delivery, high availability, and horizontal scalability. Unlike a URL shortener or a feed, a chat system is write-heavy and real-time, so the architecture must optimize for pushing data to connected users rather than only serving reads.

High-level architecture

Clients maintain persistent WebSocket connections to stateless connection/gateway servers. WebSocket is chosen over normal HTTP polling because it provides full-duplex, low-latency delivery while using TCP's reliable ordered stream. These gateway nodes authenticate users and register which user and device is connected to which node in a distributed session store such as Redis or etcd. Behind the gateways, separate services handle messaging, presence, group membership, notifications, and media. Service discovery helps these services locate each other as nodes scale up and down.

Message flow and ordering

When Alice sends a message to Bob, Alice's client includes a client-generated message ID and the conversation ID. The message service receives the message over the gateway, assigns a server-side monotonic sequence number for that conversation, and durably writes the message to a store sharded by conversation ID. Only after the write succeeds does the service return a server acknowledgement to Alice with the assigned sequence number. If Bob is online, the service consults the presence and session store to find the gateway holding Bob's connection, then pushes the message to Bob's client. Bob's client returns a delivery or read acknowledgement. If Bob is offline, the message is placed in an undelivered queue or per-user inbox so it can be delivered when Bob reconnects. Message IDs and per-conversation sequence numbers are essential for retries, idempotency, and preserving order on multi-device clients.

Presence and online status

Presence should be designed as an eventually consistent system. Clients send periodic heartbeats over their persistent connection, and the gateway refreshes a TTL-backed online flag and last-seen timestamp in Redis. If heartbeats stop or the TCP connection closes, the presence entry expires. This design favors availability and low cost over perfectly accurate presence. It is usually acceptable for online status and last-seen time to be slightly stale for a few hundred milliseconds, while message persistence and acknowledgements must be much stronger.

Storage model

Relational data such as users, groups, and memberships fits a SQL database with transactions. Message history is different: it is high-volume, append-heavy, rarely updated, and queried mostly by conversation. A good choice is to shard messages by conversation ID using consistent hashing or range-based partitioning. That keeps fetching a conversation's history local to one shard, although very active groups can create hot shards. A hybrid inbox model often works best: write the message once to the conversation log, then fan out a copy or reference to each recipient's inbox for ordinary one-on-one and small-group chats. For very large groups or broadcast channels, it is cheaper to keep the group log and let clients read from it on demand rather than duplicating a message into thousands of inboxes.

Group chat and media handling

For small and medium groups, fan-out on write is simple and fast: the messaging service persists the group message once, then enqueues delivery tasks or pushes the message to each member's connected devices. For very large groups, eager fan-out becomes expensive, so the system should switch to fan-out on read, where clients pull the latest group messages from the shared group log. Media files do not travel through the chat servers. A client uploads an image or video to object storage, receives an object URL, and sends that URL as part of the message. The receiving client downloads the object from object storage through a CDN, keeping the chat servers focused on small, real-time messages and metadata.

Multi-device and security

Each device has its own session. When a user logs in on a phone and a laptop, both register separate connections. For one-on-one messages, the system fans out to all active sessions for the recipient's user ID. If end-to-end encryption is required, clients exchange or register public keys through a key service, and messages are encrypted before they reach the server. The server still handles routing, persistence, and ordering, but it stores only ciphertext. Group encryption is more complex and typically uses a group session key distributed to all members.

Scaling and failure behavior

The stateless gateway tier scales horizontally behind load balancers because connection routing state lives in Redis or etcd, not on local disk. The message service can be partitioned by conversation or user. Asynchronous task queues and workers handle fan-out, push notifications, and media processing. Databases use replication: read replicas serve history, while the write path remains durable enough that after a server acknowledges a message, the message is not lost if a node fails. The central trade-off is to accept eventual consistency for presence and last-seen data while keeping accepted-message persistence strong and ordered. This combination of persistent WebSockets, sharded message storage, presence with TTLs, and hybrid fan-out is what allows a WhatsApp-like system to deliver billions of small messages per day with low latency and high reliability.


Design a Web Crawler

A web crawler is a system that starts from a set of seed URLs and repeatedly fetches pages, extracts links, and stores content so those pages can later be searched or analyzed. The core design challenge is balancing aggressive parallel fetching against politeness to the sites being crawled, while avoiding duplicate work and keeping the crawl scalable.

The central component is the URL frontier, which is a persistent, prioritized queue of URLs waiting to be fetched. In a real crawler this is not one simple FIFO queue. It must group URLs by domain so workers do not hit the same site too hard, and it must support priorities so important pages are fetched sooner than low-value pages. The frontier therefore acts both as a scheduler and as back pressure: when producers discover URLs faster than fetchers can drain them, queues grow, and admission control or throttling prevents memory exhaustion.

Fetch workers pull URLs from the frontier, resolve DNS, and issue HTTP requests. They must implement several correctness and politeness mechanisms. First, they fetch and cache the host's robots.txt file to avoid crawling paths that the site owner has disallowed. Second, they enforce a minimum delay between requests to the same domain and often use per-domain queues so one slow or large site cannot starve others. Third, they use timeouts, retries with exponential backoff, and per-domain rate limits so crawler traffic looks like a well-behaved client rather than an attack.

The parser stage extracts outbound links from each fetched page, normalizes them by removing fragments, resolving relative URLs, and canonicalizing hosts and paths. Because the web contains many duplicate and near-duplicate URLs, the crawler deduplicates aggressively before a URL is allowed back into the frontier. A Bloom filter provides a memory-efficient test for whether a URL was seen before, at the cost of occasional false positives; the frontier and storage layer maintain the authoritative seen set. Some crawlers also use content fingerprints or shingling to avoid repeatedly indexing identical content under different URLs.

Fetched pages and metadata are stored in a decoupled way. Raw HTML or rendered content usually goes into a distributed object store, while typed metadata such as URL, status, title, canonical link, and fetch time goes into a database that supports efficient queries for indexing and future recrawls. Offline or streaming jobs build the search index from this stored corpus.

For scale, the crawler is built from stateless workers. The frontier can be backed by a distributed queue or log such as Kafka, allowing many fetchers and parsers to consume work independently. The seen-set and robots cache are shared services, and the whole crawl can be partitioned by domain or by URL hash to spread work across machines. Coordination across machines is necessary because two workers in different places must not bypass the global per-domain politeness limit.

The main trade-offs are freshness versus load and completeness versus cost. A fast crawler can overwhelm sites and violate accepted crawling norms, while an overly polite crawler may never keep up with a large or changing web. The design therefore prioritizes robustness, deduplication, respect for robots.txt, and domain-level fairness, and it uses queues, caches, and distributed workers to scale out rather than relying on a single fast machine.


Requirements and Scale for WhatsApp-like Chat

To size the system, assume 500 million daily active users, each sending around 30 messages per day on average. That works out to 15 billion messages per day total, or roughly 174,000 messages per second average. Since chat traffic is bursty and peaks at maybe 3x average, plan for a peak of about 500,000 messages per second. Storage: a typical message with metadata and index overhead is about 200 bytes; 15 billion messages per day means 3 TB of raw message data per day, or about 1.1 PB per year. If we retain history for five years and include attachments (photos, videos), the media storage will be much larger and should be kept in object storage, while the message table itself remains relatively small. These numbers justify choices like sharding by conversation ID, using a durable message log, and separating presence (cheap, eventually consistent) from message persistence (expensive, strongly consistent).