MUHAMMAD
FARI
MADYAN
[ Press ESC or Click to Skip ]

System Design Roadmap

1

Part 1: Networking Basics – Packets, TCP/UDP, TLS & REST

2

Part 2: Core Building Blocks – Load Balancers, CDN, Caching, API Gateway, DBs & Kafka

3

Part 3: Distributed Patterns & Concepts – CAP Theorem, Hashing, Sharding & Sagas

4

Part 4: Real-World Interview Systems – URL Shortener, Chat, Food Delivery & Notification

This article is available in Indonesian

🇮🇩 Baca dalam Bahasa Indonesia
🇬🇧 English📚 System Design Roadmap

System Design Roadmap Part 3: Distributed Patterns & Concepts – CAP Theorem, Hashing, Sharding & Sagas

Deep-dive into distributed systems theory and production patterns: CAP and PACELC theorems, consistent hashing rings, horizontal database sharding, idempotent APIs, Transactional Outbox, Saga orchestration, and CQRS.

Muhammad Fari MadyanAuthor

12 min read

·

Oct 7, 2026


Introduction: Moving Beyond Simple Lego Blocks

In Part 2, we laid out the physical components of scale: load balancers, caches, databases, and message queues.

However, junior and mid-level engineers treat these components like straightforward Lego blocks, assuming that connecting Kafka to Postgres will just magically work. Senior and Principal engineers know the painful reality: in distributed systems, anything that can go wrong will go wrong.

Network cables get severed by undersea earthquakes. Nodes crash in the middle of executing a bank transfer. Duplicate HTTP requests arrive out of order. Database disks fill up.

To pass senior system design interviews, you must demonstrate a mastery of distributed systems theory and the battle-tested architectural patterns engineered to survive chaos.

┌─────────────────────────────────────────────────────────────────────────────┐ │ THE SPECTRUM OF DISTRIBUTED SYSTEM PATTERNS │ └─────────────────────────────────────────────────────────────────────────────┘ [ Theoretical Limits ] ──► CAP Theorem & PACELC Trade-offs │ [ Data Partitioning ] ──► Consistent Hashing & Horizontal Sharding │ [ Fault Tolerance ] ──► Idempotency Keys & Deduplication │ [ Microservice Sync ] ──► Transactional Outbox, Saga Pattern & CQRS │ [ Traffic Defense ] ──► Distributed Rate Limiting & Token Buckets

1. The CAP Theorem & The PACELC Reality

Formulated by Eric Brewer in 2000, the CAP Theorem states that in any asynchronous distributed data store, you can only simultaneously guarantee two out of the following three properties:

                 Consistency (C)
                    ▲       ▲
                   /         \
                  /   RDBMS   \
                 /  (Single)   \
                /               \
Availability (A) ════════════════ Partition Tolerance (P)
        (DynamoDB, Cassandra)       (HBase, MongoDB, Spanner)
  1. Consistency (C): Every read receives the most recent write or an error (Linearizability).
  2. Availability (A): Every non-failing node returns a non-error response for every request (without guaranteeing it is the latest write).
  3. Partition Tolerance (P): The system continues to operate despite arbitrary network message drops or delays between nodes.

The Interview Truth: "CA" Does Not Exist in Distributed Systems

Network partitions are an inevitable physical reality of distributed infrastructure (cables break, switches fail, packet loss occurs). Therefore, Partition Tolerance (P) is mandatory.

Your real choice is between CP and AP:

  • CP Systems (Consistency over Availability): If Node 1 cannot talk to Node 2, the system rejects writes or returns errors rather than risk serving stale or divergent data. Examples: Google Cloud Spanner, Apache HBase, MongoDB (majority write concern), banking core ledgers.
  • AP Systems (Availability over Consistency): If network partitions occur, both nodes continue accepting writes independently. The system accepts temporary data divergence and resolves conflicts later via eventual consistency. Examples: Apache Cassandra, Amazon DynamoDB, Couchbase, DNS.

Beyond CAP: The PACELC Theorem

CAP only describes system behavior during a network partition. What happens during normal 99.9% uptime?

Daniel Abadi's PACELC Theorem completes the picture:

If there is a Partition, trade off Availability vs Consistency; Else, trade off Latency vs Consistency.

  • PA/EL: If partitioned, choose Availability; Else, prioritize low Latency over strong Consistency (e.g., DynamoDB, Cassandra).
  • PC/EC: If partitioned, choose Consistency; Else, prioritize strong Consistency over Latency (e.g., PostgreSQL with synchronous replication, Bigtable).

2. Consistent Hashing: Minimizing Key Churn in Dynamic Clusters

When distributing millions of cache keys across $N$ Redis or Memcached servers, naive modulus hashing seems tempting:

$$\text{Server Index} = \text{hash}(\text{key}) \pmod N$$

The Catastrophe of Modulus Hashing

If you have 4 servers ($N=4$) and add 1 new server ($N=5$), the denominator changes for every single key. Roughly 80% of all cached keys instantly map to different servers! This causes a massive, instantaneous cache miss across your entire infrastructure, flooding the primary database and triggering a catastrophic outage.

The Consistent Hashing Ring

Consistent hashing maps both servers and data keys onto a circular 32-bit Hash Ring ($0$ to $2^{32} - 1$):

                                  [ Server A ] (0)
                                     /       \
                                    /         \
                      [ Key 1 ] ───►           \
                                  /             \
                   [ Server C ]                  [ Server B ]
                   (2^32 * 0.66)                (2^32 * 0.33)
                                  \             /
                                   \           /
                     [ Key 2 ] ─────►         /
                                     \       /
  1. Hash the server identifiers (e.g., Server-A-IP) to positions on the ring.
  2. Hash incoming keys (e.g., user_49102) to positions on the same ring.
  3. To find which server stores a key, walk clockwise along the ring until you encounter the first server.

What Happens When a Server Is Added or Removed?

When a new server is inserted between Server A and Server B, only the keys between Server A and the new server are relocated. All other keys on the rest of the ring remain untouched! When $N$ nodes exist, adding or removing a node migrates only $1/N$ of the total keys.

Virtual Nodes (Vnodes) Solve Hotspots

In a basic ring, server hashes may land unevenly, creating severe hotspots (one server holding 60% of traffic).

Solution: Assign each physical server hundreds of Virtual Nodes (vnodes) scattered randomly across the ring (e.g., Server-A#1, Server-A#2, ..., Server-A#200). This homogenizes data distribution to within 1-2% variance and enables weighted capacity (giving a 64GB server 200 vnodes and a 32GB server 100 vnodes).


3. Database Sharding: Scaling Writes Horizontally

When a relational database table exceeds tens of millions of rows, indexes no longer fit in RAM, and write IOPS saturate storage controllers. Sharding (Horizontal Partitioning) splits a single logical table across multiple autonomous database instances.

┌─────────────────────────────────────────────────────────────────────────────┐ │ DATABASE SHARDING TOPOLOGY │ └─────────────────────────────────────────────────────────────────────────────┘ [ Application ] │ ▼ [ Sharding Routing Layer ] │ ┌───────────────────────┼───────────────────────┐ ▼ ▼ ▼ [ Shard 1 ] [ Shard 2 ] [ Shard 3 ] (Users 1 - 10M) (Users 10M - 20M) (Users 20M - 30M)

Sharding Strategies

  1. Range-Based Sharding:
    • Routes data based on discrete value ranges (e.g., Shard 1 holds User IDs 1 to 1,000,000; Shard 2 holds 1,000,001 to 2,000,000).
    • Risk: Severe write hotspots. Since auto-increment IDs grow sequentially, 100% of new writes hit the highest shard while older shards sit idle.
  2. Directory-Based (Lookup) Sharding:
    • Maintains a centralized lookup table mapping Entity ID -> Shard ID.
    • Advantage: Ultimate flexibility; individual shards can be rebalanced dynamically.
    • Risk: The lookup service becomes a single point of failure and adds latency to every query.
  3. Hash-Based Sharding (Industry Standard):
    • Routes rows using $\text{hash}(\text{ShardKey}) \pmod{\text{Number of Shards}}$.
    • Advantage: Uniform write distribution across all nodes.
    • Trade-off: Re-sharding requires complex background data migration pipelines.

The Sharding Nightmares You Must Address in Interviews:

  • Cross-Shard Joins: Executing JOIN queries across tables located on different physical database servers is prohibitively slow. Solution: Denormalize data, duplicate lookup tables across all shards, or execute application-level joins.
  • Distributed Transactions: Modifying data across multiple shards requires two-phase commits (2PC), which introduces massive lock contention and latency. Solution: Design the Shard Key so that related entities live on the exact same shard (e.g., shard orders by user_id so all of a user's orders reside together).

4. Idempotency: Defeating the Unreliable Network

In distributed architectures, networks drop packets. When a client sends a payment request and experiences an HTTP 504 Gateway Timeout, the client cannot know whether the server failed before processing, or if the server succeeded but the response packet was lost on the wire.

If the client naively retries the request, the user may be charged twice.

Client Server │ │ │ ── POST /v1/payments (Idempotency-Key: abc-123) ────> │ │ │ 1. Checks Redis: Not found │ │ 2. Sets key in Redis: "PROCESSING" │ │ 3. Executes Bank Charge │ │ 4. Updates Redis: "SUCCESS", payload │ X (Response Dropped by Network Flap!) │ │ │ │ ── RETRY POST /v1/payments (Key: abc-123) ──────────> │ │ │ 5. Checks Redis: Found "SUCCESS" │ <── HTTP 200 OK (Cached Original Result) ──────────── │ 6. Returns immediately!

Implementing Enterprise Idempotency:

  1. The client generates a unique UUID Idempotency-Key header with every mutating request.
  2. The server uses an atomic check-and-set in Redis (SET key "PROCESSING" NX EX 120).
    • If the key already exists with status PROCESSING, return HTTP 409 Conflict (request currently in flight).
    • If the key already exists with status COMPLETED, return the cached JSON response immediately without touching the business logic.
  3. The database enforces a UNIQUE(idempotency_key) constraint on the transaction table to prevent race conditions at the storage tier.

5. Distributed Transaction Patterns: Outbox, Saga & CQRS

In a microservices architecture, each service owns its private database. A monolithic ACID transaction spanning across services is impossible. How do we ensure consistency?

Pattern 1: The Transactional Outbox Pattern (Eliminating the Dual-Write Problem)

THE DUAL-WRITE FLAW:
1. OrderService writes Order to PostgreSQL database.
2. OrderService tries to emit event to Apache Kafka.
3. Kafka cluster is unreachable!
Result: Database has the Order, but Kafka never gets the event. Systems desynchronize!
┌─────────────────────────────────────────────────────────────────────────────┐ │ TRANSACTIONAL OUTBOX PATTERN │ └─────────────────────────────────────────────────────────────────────────────┘ [ Order Service ] │ │ 1. Atomic Local SQL Transaction: │ - INSERT INTO orders (...) │ - INSERT INTO outbox_table (event_payload, status) ▼ [(PostgreSQL Database)] │ │ 2. Transaction Log Tailer / CDC (e.g., Debezium) │ Reads PostgreSQL WAL (Write-Ahead Log) ▼ [ Debezium / Poller ] │ │ 3. Publishes guaranteed message ▼ [ Apache Kafka ]

By persisting the event into an outbox table inside the exact same local ACID transaction as the order record, we guarantee that if the database write succeeds, the event is guaranteed to be recorded. A background CDC engine (Debezium reading Postgres WAL) or poller streams the record reliably into Kafka.


Pattern 2: The Saga Pattern (Distributed Transactions Without 2PC)

A Saga represents a long-running business transaction split into a sequence of local microservice transactions. If any step fails, the Saga executes Compensating Transactions in reverse order to roll back state.

Happy Path:
[ Create Pending Order ] ──► [ Reserve Inventory ] ──► [ Process Payment ] ──► [ Order Confirmed ]

Failure & Compensation:
[ Create Pending Order ] ──► [ Reserve Inventory ] ──► [ Payment Fails! ]
                                       │
                                       ▼ Compensating Step
                             [ Unreserve Inventory ] ──► [ Mark Order Cancelled ]

Choreography vs. Orchestration:

  • Choreography (Event-Driven): Services listen to Kafka events and trigger their next steps autonomously. Best for simple workflows (2–4 services). Danger: Hard to track and debug when complexity grows ("spaghetti events").
  • Orchestration (Central Coordinator): A dedicated state machine orchestrator (e.g., Temporal, AWS Step Functions, custom orchestrator) explicitly issues commands to services and tracks overall Saga status. Best for complex, mission-critical business flows.

Pattern 3: CQRS (Command Query Responsibility Segregation)

In high-scale systems, read workloads and write workloads have radically different performance, latency, and consistency requirements (often 99% reads, 1% writes).

┌─────────────────────────────────────────────────────────────────────────────┐ │ CQRS ARCHITECTURE PATTERN │ └─────────────────────────────────────────────────────────────────────────────┘ [ Client Application ] / \ Write (Command)/ \ Read (Query) / \ ▼ ▼ [ Command Service ] [ Query Service ] │ │ ▼ ▼ [ Primary RDBMS ] [ Read Replicas / ] (Normalized SQL) [ Elasticsearch ] │ (Denormalized JSON) │ ▲ │ 1. Emits Events │ └──────► [ Kafka ] ────┘ 2. Syncs Projections
  • Command Model: Handles writes, enforces complex business validation invariants, and writes to a normalized relational database (PostgreSQL).
  • Query Model: Emits domain events via Kafka to asynchronous projection workers that maintain read-optimized, pre-computed views in Elasticsearch, MongoDB, or Redis. Queries execute in single-digit milliseconds without joins!

6. Distributed Rate Limiting: Protecting APIs from Traffic Spikes

Rate limiting throttles client requests to protect backends from abusive scrapers, credential stuffing, and sudden traffic spikes.

┌──────────────────────────────────────┬──────────────────────────────────────┐ │ Algorithm │ Pros vs. Cons │ ├──────────────────────────────────────┼──────────────────────────────────────┤ │ 1. Token Bucket │ Memory efficient; allows bursts. │ │ 2. Leaky Bucket │ Smooths traffic spikes into steady │ │ │ trickle; drops requests if queue full│ │ 3. Fixed Window Counter │ Simple; vulnerable to 2x boundary │ │ │ traffic bursts at window edges. │ │ 4. Sliding Window Log │ 100% accurate; high memory footprint │ │ │ (stores every timestamp in ZSET). │ │ 5. Sliding Window Counter │ Memory lightweight; smoothly blends │ │ │ adjacent windows with <0.05% error. │ └──────────────────────────────────────┴──────────────────────────────────────┘

Production Implementation: Redis Lua Sliding Window Counter

In a distributed multi-server cluster, rate limiting must be atomic. Executing multiple round-trip calls to Redis introduces race conditions. We solve this by executing an atomic Lua Script inside Redis:

LUA

-- KEYS[1]: Rate limit key (e.g., "ratelimit:user_123:minute")
-- ARGV[1]: Max requests allowed (e.g., 100)
-- ARGV[2]: Current timestamp in seconds
-- ARGV[3]: Window size in seconds (e.g., 60)

local current = redis.call('GET', KEYS[1])
if current and tonumber(current) >= tonumber(ARGV[1]) then
    return 0 -- Rate limit exceeded! Return HTTP 429
else
    local count = redis.call('INCR', KEYS[1])
    if tonumber(count) == 1 then
        redis.call('EXPIRE', KEYS[1], ARGV[3])
    end
    return 1 -- Request allowed
end

By running the script directly within Redis, the execution is single-threaded and guaranteed atomic across thousands of concurrent API Gateway workers.


Architectural Synthesis: The Resilience Matrix

ChallengeDistributed Solution
Undersea cable failure cuts cross-datacenter syncDesign for AP (eventual consistency) or reject writes gracefully (CP).
Adding new cache nodes triggers mass evictionConsistent Hashing with virtual nodes to limit churn to $1/N$.
Relational table exceeds 500M recordsHash-Based Database Sharding keyed by tenant or user ID.
Network timeout on credit card paymentClient passes Idempotency-Key, server verifies atomic state in Redis.
Microservice crash halfway through order flowSaga Pattern Orchestrator invokes compensating reversal steps.
DDoS attack or runaway client scriptSliding Window Rate Limiter at API Gateway returning HTTP 429 Too Many Requests.

What's Next in Part 4?

Now that we have conquered networking fundamentals, core components, and distributed architectural patterns, it is time to put everything to the test under real interview conditions.

In the final chapter, Part 4: Real-World Interview Systems, we will design 5 complete end-to-end production systems from scratch:

  1. Design a URL Shortener (TinyURL)
  2. Design a Real-Time Chat System (WhatsApp)
  3. Design a Food Delivery App (Zomato / Gojek)
  4. Design a Scalable Notification System
  5. Design a Distributed Rate Limiter

Continue Reading

Previous article

← Previous Article

System Design Roadmap Part 2: Core Building Blocks – Load Balancers, CDN, Caching, API Gateway, DBs & Kafka

Next Article →

System Design Roadmap Part 4: Real-World Interview Systems – URL Shortener, Chat, Food Delivery & Notification

Next article