A production-ready system design cheat sheet provides immediate access to hardware latency numbers, capacity estimation formulas, storage engine trade-offs, and distributed consensus primitives. When designing multi-region architectures or answering architectural prompts during senior and staff technical rounds, relying on intuition produces vague, non-viable topologies. Precision engineering demands exact numbers, deterministic sizing, and structural knowledge of where networks and disks fail.
Distributed systems at hyperscale encounter failure as a baseline state rather than an anomaly. Networks partition, tail latencies cascade across microservice meshes, disks saturate on write-heavy compaction cycles, and caches drop to near-zero hit rates during hotkey events. A candidate or lead engineer must transition rapidly from broad system requirements to concrete bandwidth budgets, throughput thresholds, and replication guarantees.
This technical reference serves as an operational blueprint for capacity planning, component selection, failure domain isolation, and resilient API boundary construction under extreme scale.
Hardware Latency and Back-of-the-Envelope Estimation Constants
Reliable capacity planning begins with memorized hardware physics. You cannot design a durable message log or a multi-tier cache without understanding the four-to-six order-of-magnitude performance cliff separating CPU registers, main memory, local non-volatile storage, and physical network hops. Incorporating this system design cheat sheet reference into your calculations guarantees that latency budgets remain grounded in production reality.
Latency Numbers Every Systems Architect Must Know
Peter Norvig and Jeff Dean established foundational latency numbers that every infrastructure engineer must internalize. The numbers below reflect modern datacenter hardware, including PCIe Gen 5 NVMe drives, DDR5 memory, and high-performance top-of-rack switches.
| Hardware Operation | Approximate Latency (2026) | Scaled Metric (1 ns = 1 sec) | Engineering Implication |
|---|---|---|---|
| L1 Cache Reference | 0.5 to 1 ns | 1 second | Extremely fast; instruction-level optimization territory. |
| L2 Cache Reference | 3 to 4 ns | 3.5 seconds | Fast on-die memory access. |
| L3 Cache Reference | 10 to 20 ns | 15 seconds | Shared CPU cache across cores. |
| Main Memory (DDR5 RAM) Read | 50 to 100 ns | 1.5 minutes | Primary baseline for in-memory databases like Redis. |
| NVMe SSD Random Read | 10 to 50 microseconds | 11.5 hours | Random I/O on NVMe is orders of magnitude faster than spinning disk, but still 500x slower than RAM. |
| Sequential SSD Read (1 MB) | 50 to 100 microseconds | 1.1 days | Sequential scans leverage controller prefetching; preferred for log-structured merge trees. |
| Intra-Datacenter Network RTT | 200 to 500 microseconds | 4.5 days | Direct RPC limit; microservices traversing availability zones compound tail latencies here. |
| Cross-Region Network RTT (e.g. US-East to US-West) | 60 to 70 ms | 2 years | Synchronous two-phase commit over this distance destroys write throughput. |
| Cross-Continent Network RTT (e.g. US to Europe) | 100 to 150 ms | 4.5 years | Requires edge termination and asynchronous replication. |
| Cold Disk Seek (Mechanical HDD) | 4 to 10 ms | 4 months | Unacceptable for direct user-facing queries; acceptable only for deep archival or bulk sequential cold storage. |
[ L1/L2/L3 Cache: ~1-20 ns ] ──> [ RAM: ~100 ns ] ──> [ NVMe Read: ~20 µs ] ──> [ Intra-DC RTT: ~500 µs ] ──> [ WAN RTT: ~100 ms ]
└────── In-Memory Compute ──────┘ └──────── Local I/O ────────┘ └──────── Network Boundaries ────────┘
Powers of Two and Storage Sizing Conversions
Capacity math requires instant conversions between bit powers, byte notations, and query counters without a calculator:
2^10 = 1,024 = ~1 Thousand (1 KB)2^20 = 1,048,576 = ~1 Million (1 MB)2^30 = 1,073,741,824 = ~1 Billion (1 GB)2^40 = 1,099,511,627,776 = ~1 Trillion (1 TB)2^50 = ~1 Quadrillion (1 PB)
Standard Estimation Formulas
When computing queries per second (QPS), bandwidth, and disk consumption, apply standard temporal and margin conversions:
- Seconds Per Day: Exactly 86,400 seconds. For rapid calculation, round up to
100,000or86,400 ≈ 8.64 × 10^4. Using 100,000 provides an automatic 15% safety buffer for daily volume. - Read QPS:
(Daily Active Users × Reads per User per Day) / 86,400. - Write QPS:
(Daily Active Users × Writes per User per Day) / 86,400. - Peak QPS: Apply a burst multiplier based on the domain. Standard enterprise platforms experience a
2x to 3xpeak, while ticketing platforms, flash sales, and live sports systems peak between5x and 10xbaseline QPS. - Ingress Bandwidth:
Write QPS × Payload Size per Write. - Egress Bandwidth:
Read QPS × Payload Size per Read. - Storage Growth (3 Years):
Write QPS × Payload Size × 86,400 × 365 × 3 × Replication Factor (e.g. 3) × Index/Metadata Overhead (e.g. 1.4).
Capacity Calculation Rule: Always isolate read footprints from write footprints. Never sum them into a single aggregate bandwidth metric. Reads depend heavily on cache hit ratios, edge routing, and replica pools, whereas writes dictate disk compaction, journal serialization, consensus quorum traffic, and permanent storage scaling.
The 45-Minute System Design Interview Execution Framework
A high-scoring interview requires strict time allocation. System design discussions collapse when candidates dive into caching or database schemas before establishing unambiguous functional boundaries and throughput targets. Use this structured system design interview cheat sheet process to navigate a 45-minute architectural interview.
- Phase 1: Scope, Clarify, and Estimate (0 to 5 Minutes)
Establish hard system boundaries. Inquire about user demographics, read-to-write ratios, latency thresholds (e.g. p99 under 100ms), and durability requirements. Calculate peak read/write QPS, network ingress/egress, and five-year storage growth using the estimation constants above. Agree on two to three core functional user journeys and mark out-of-scope features explicitly. - Phase 2: High-Level Architecture and Data Flow (5 to 15 Minutes)
Draft the primary components: Client, Edge CDN/API Gateway, Stateless Application Services, Message Brokers, and Storage. Define the core data entities and establish primary HTTP or gRPC API contracts (e.g.POST /v1/orderswith idempotency keys). Ensure every component directly serves one of the functional requirements established in Phase 1. - Phase 3: Deep-Dive Component Design and Bottlenecks (15 to 35 Minutes)
Drill down into the primary technical complexity. If the prompt is a distributed messaging system, design the log segment manager, write-ahead log (WAL) flushing, and consumer group offset management. If designing a feed platform, solve fan-out mechanics (fan-out-on-write vs fan-out-on-read). Address database partitioning, shard keys, distributed locking, and multi-layer caching patterns. - Phase 4: Resilience, Failure Modes, and Observability (35 to 45 Minutes)
Audit the architecture for single points of failure (SPOFs). Walk through concrete disaster scenarios: primary node crash, split-brain condition, cascading service failures, hot partition overload, and cache thundering herds. Describe circuit breaker trip conditions, fallback strategies, distributed tracing with OpenTelemetry, and cross-region disaster recovery (RPO and RTO).
Architectural Checklist Before Whiteboarding
- Have you clarified whether consistency or availability takes precedence during an AZ partition?
- Are the ingress payloads secured with idempotency keys to handle client-side retries?
- Is read traffic decoupled from the transactional primary database via replicas or read-through caches?
- Have you identified the explicit sharding key and evaluated its distribution across hash rings?
- Are asynchronous batch operations isolated from synchronous critical-path request loops?
+-----------------------------------------------------------------------------------------+
| TIME-BOX EXECUTION MATRIX |
+-----------------------------+-----------------------------------+-----------------------+
| Time Window | Stage | Primary Deliverable |
+-----------------------------+-----------------------------------+-----------------------+
| 00:00 - 05:00 (5 mins) | Functional & Scale Scoping | QPS, SLA, Storage Math|
| 05:00 - 15:00 (10 mins) | High-Level Topology & APIs | End-to-End Diagram |
| 15:00 - 35:00 (20 mins) | Deep-Dive Critical Component | Data Structures/Shards|
| 35:00 - 45:00 (10 mins) | Edge Cases, Failures, Metrics | MTTR, SPOF Elimination|
+-----------------------------+-----------------------------------+-----------------------+
Storage Engine Decision Matrix: SQL, NoSQL, and Partitioning
Choosing a storage engine without analyzing disk access patterns, write-amplification factors, and index traversal costs is a major architectural mistake. This system design cheatsheet segment breaks down storage paradigms, underlying disk structures, and partitioning algorithms.
Storage Engine Architecture Comparison
| Storage Category | Underlying Data Structure | Optimized Access Pattern | ACID Guarantees | Common Production Solutions |
|---|---|---|---|---|
| Relational (RDBMS) | B+ Tree (Page-based, 4KB-16KB pages) | Point lookups, range queries, multi-row ACID joins | Full ACID (Serializable / Snapshot Isolation) | PostgreSQL, MySQL (InnoDB), CockroachDB |
| Key-Value Store | In-Memory Hash Table / SkipList + AOF/RDB | O(1) point reads/writes by exact key | Single-key atomic; distributed multi-key via locking | Redis, KeyDB, AWS DynamoDB |
| Wide-Column / LSM-Tree | Log-Structured Merge Tree (SSTables + MemTable) | High-throughput append-only writes, sparse column lookups | Row-level atomicity; eventual consistency tunable | Apache Cassandra, ScyllaDB, RocksDB |
| Document Store | B-Tree or JSON-optimized variant (WiredTiger) | Polymorphic documents, nested data structures | Document-level ACID, multi-document via replica sets | MongoDB, Couchbase |
| Time-Series Engine | Delta-of-Delta, Gorilla, Sparse Indexing | High-frequency time-stamped telemetry writes, aggregation | Append-focused; limited transactional support | TimescaleDB, ClickHouse, InfluxDB |
Data Partitioning Schemes
When storage volume or write QPS surpasses the physical limits of a single server, data must be sharded across a cluster. The choice of partitioning algorithm dictates balancing, latency, and resharding overhead:
- Range-Based Partitioning: Keys are partitioned in continuous sorted sequences (e.g. A-C, D-F). While range queries are natural and efficient, this approach produces intense write hot-spotting when keys are monotonically increasing, such as auto-incrementing IDs or timestamps.
- Hash-Based Partitioning: Keys pass through a uniform hashing function (e.g.
MurmurHash3(key) mod N). This scatters data evenly across storage nodes. The fatal flaw is resizing: adding or removing a node changesN, triggering a migration of almost all stored keys across the network. - Consistent Hashing: Nodes and keys map to a shared 360-degree integer ring (typically
[0, 2^32 - 1]). A key routes to the first node encountered moving clockwise. When adding or removing a node, onlyK/Nkeys must migrate, whereKis the total key count andNis the total node count. - Virtual Nodes (Vnodes): To solve the uneven distribution problem in consistent hashing, where physical nodes hold disproportionate segments of the ring, each physical node is assigned 100 to 256 virtual positions on the ring. This balances memory consumption, distributes network load uniformly, and prevents hot spot creation during node additions or removals.
LSM-Tree vs B+ Tree Rule: Use B+ Tree engines (PostgreSQL) when read performance on arbitrary ranges and immediate consistency are strict requirements. Use LSM-Tree engines (Cassandra, RocksDB) when write throughput dominates read traffic, because LSM trees convert costly random disk writes into sequential writes in a memory buffer (MemTable) before flushing immutable SSTables to disk.
Distributed Consensus and CAP vs PACELC Trade-Offs
Every distributed architecture must handle network partitions, transient hardware failures, and replication lag. Treating network communication as instantaneous or completely reliable triggers split-brain conditions and silent data corruption.
CAP Theorem vs PACELC Theorem
The CAP theorem states that in the event of a network Partition (P), an architect must choose between Consistency (C) and Availability (A). However, network partitions are rare operational events. Eric Brewer’s CAP model fails to describe how systems behave during the 99.9% of normal operational runtime.
The PACELC theorem resolves this limitation: If there is a Partition (P), how does your system choose between Availability (A) and Consistency (C)? Else (E), when the system runs normally without partitions, how do you trade off Latency (L) versus Consistency (C)?
| System | PACELC Classification | Partition Behavior (P/A or P/C) | Normal Runtime Behavior (E/L or E/C) | Design Implications |
|---|---|---|---|---|
| Apache Cassandra / DynamoDB (Default) | PA/EL | Chooses Availability; continues local writes on separated nodes | Chooses Latency; replicates asynchronously, yielding eventual consistency | Ultra-low tail latency; susceptible to stale reads and conflicting concurrent writes. |
| Apache HBase / Bigtable | PC/EC | Chooses Consistency; rejects writes if master or region server is unconfirmed | Chooses Consistency; forces client reads through primary to avoid stale data | Strong consistency; read/write operations stall if a node fails failover checks. |
| MongoDB (with Majority Write Concern) | PC/EC | Chooses Consistency; blocks writes until majority acknowledgment succeeds | Chooses Consistency; primary handles operations by default | Prevents dirty reads and rollbacks; incurs network round-trip delays on writes. |
| Amazon DynamoDB (Strong Consistent Reads) | PA/EC | Chooses Availability during partition events | Chooses Consistency when explicitly requested via SDK query flags | Allows dynamic query-level trade-offs between cost, latency, and read freshness. |
Distributed Consensus Protocols
When multiple nodes must agree on a state machine transition or a transaction commit sequence, distributed consensus algorithms coordinate updates:
- Raft: Deconstructs consensus into explicit sub-problems: Leader Election, Log Replication, and Safety. A cluster requires a strict quorum of
(N/2) + 1nodes to elect a leader and commit a log entry. Raft powers etcd and CockroachDB due to its clear state transitions and formal verifiability. - Multi-Paxos: The historic production standard (used in Google Spanner and Chubby). Paxos separates state proposers, acceptors, and learners. While theoretically equivalent to Raft in performance, its implementation complexity often leads to subtle bugs in edge-case leader failovers.
- Quorum Intersection (Dynamo-style Leaderless Systems): Utilizes configurable client-side parameters:
N(number of replicas),W(number of replicas that must acknowledge a write), andR(number of replicas that must respond to a read).
The Strict Quorum Inequality: Strong consistency requires
W + R > N. IfN = 3, configuringW = 2andR = 2guarantees that the read set and the write set intersect at least on one node containing the latest monotonic timestamp or vector clock. IfW + R <= N, the system operates in an eventual consistency mode with lower read and write latencies.
Caching Topologies and Cache Invalidation Failure Modes
Caching introduces mutable state divergence between primary storage and memory layers. Poorly implemented caching strategies cause thundering herds, cold-start query cascading, and silent data corruption.
Cache Topologies and Access Patterns
CACHE-ASIDE (Lazy Loading) WRITE-THROUGH WRITE-BEHIND (Write-Back)
+--------+ 1. Read Miss +-------+ +--------+ 1. Write +-------+ +--------+ 1. Write +-------+
| Client | ---------------> | Cache | | Client | ---------> | Cache | | Client | ---------> | Cache |
+--------+ +-------+ +--------+ +-------+ +--------+ +-------+
| 2. Fetch DB ^ | | | |
v | | 2. Sync Write | | 2. Async Queue |
+--------+ ---------------------+ v | v |
| DB | 3. Populate Cache +--------+ <--------------+ +--------+ <--------------+
+--------+ | DB | | DB | (Batch Flush)
+--------+ +--------+
- Cache-Aside (Lazy Loading): The application coordinates reads and writes directly. On a read miss, the application loads the record from the database, writes it to the cache, and returns it. Writes update the database directly and invalidate (delete) the cached entry. Advantage: Resilient to cache crashes. Disadvantage: Triple network trip on cold misses.
- Write-Through: The application writes directly to the cache. The cache synchronously writes the data to the underlying database within the same execution path. Advantage: Read misses are minimal; data is fresh. Disadvantage: High write latency due to synchronous double-writes.
- Write-Behind (Write-Back): The application writes to the cache, which acknowledges the client immediately. The cache then asynchronously flushes dirty blocks or records to the database in micro-batches. Advantage: Maximum write performance; absorbs high write bursts. Disadvantage: Data loss occurs if the cache node crashes before flushing dirty writes to durable storage.
Mitigating Catastrophic Cache Failure Modes
Distributed caches suffer from three classic failure modes that can cascade across your architecture:
- Cache Penetration: Requests query keys that do not exist in either the cache or the database (e.g. malicious requests scanning random UUIDs). These queries bypass the cache entirely and hit the database directly, causing CPU and I/O saturation. Mitigation: Deploy a Bloom filter in front of the cache to quickly reject nonexistent keys, or store null objects in the cache with a short TTL (e.g. 30 seconds).
- Cache Breakdown (Hotspot Invalidation): A single high-traffic key (such as a breaking news item or flash sale listing) expires while receiving thousands of concurrent requests. Every concurrent thread encounters a cache miss and hits the database simultaneously. Mitigation: Use distributed mutexes to allow only one thread to regenerate the key, or implement probabilistic early recomputation (XFetch).
- Cache Avalanche: A large collection of cached items share the same TTL and expire at the exact same moment, or an entire cache node crashes, redirecting immense read traffic to the primary database. Mitigation: Add uniform random jitter to all TTL values (e.g.
TTL = Base_TTL + rand(0, 300)seconds) and deploy multi-master distributed cache clusters with automated failover.
Production Implementation: Distributed Mutex Against Cache Stampedes
The Go implementation below demonstrates how to prevent the thundering herd problem using a Redis-backed distributed lock with a double-check pattern.
package cache
import (
"context"
"errors"
"math/rand"
"time"
"github.com/redis/go-redis/v9"
)
type ResilientCache struct {
rdb *redis.Client
fetchFromDB func(ctx context.Context, key string) (string, error)
}
func (c *ResilientCache) GetWithStampedeProtection(ctx context.Context, key string, baseTTL time.Duration) (string, error) {
// 1. Initial Cache Lookup
val, err:= c.rdb.Get(ctx, key).Result()
if err == nil {
return val, nil
}
if!errors.Is(err, redis.Nil) {
return "", err // Real redis network/system error
}
// 2. Cache Miss: Acquire Distributed Lock to prevent Thundering Herd
lockKey:= "lock:" + key
lockAcquired, err:= c.rdb.SetNX(ctx, lockKey, "1", 5*time.Second).Result()
if err!= nil {
return "", err
}
if!lockAcquired {
// Another worker is actively populating the cache; back off and retry
time.Sleep(50 * time.Millisecond)
return c.GetWithStampedeProtection(ctx, key, baseTTL)
}
// Ensure lock is released after fetching
defer c.rdb.Del(ctx, lockKey)
// 3. Double-check cache in case another worker populated it while acquiring lock
val, err = c.rdb.Get(ctx, key).Result()
if err == nil {
return val, nil
}
// 4. Query Database (Protected: exactly one worker executes this)
data, err:= c.fetchFromDB(ctx, key)
if err!= nil {
// Cache penetration protection: Cache empty string with short TTL
c.rdb.Set(ctx, key, "", 30*time.Second)
return "", err
}
// 5. Populate Cache with Jittered TTL to prevent Cache Avalanche
jitter:= time.Duration(rand.Intn(120)) * time.Second
totalTTL:= baseTTL + jitter
c.rdb.Set(ctx, key, data, totalTTL)
return data, nil
}
Eviction Policy Selection: For general-purpose read caches, select
allkeys-lru(Least Recently Used) orvolatile-lru. For workloads with recurring periodic batch jobs that can flush out useful keys, chooseallkeys-lfu(Least Frequently Used) to protect against cold-scan pollution.
Edge Resiliency: Rate Limiting, Circuit Breakers, and Backpressure
A distributed system that lacks edge defenses will experience cascading failures when downstream microservices degrade. Without backpressure, incoming requests queue up, consuming thread pools, sockets, and memory until the entire system crashes.
Rate Limiting Algorithms Breakdown
| Algorithm | Memory Footprint | CPU Overhead | Burst Handling | Failure Vulnerability |
|---|---|---|---|---|
| Token Bucket | Minimal: O(1) per user (tokens, timestamp) | Low: Simple mathematical token top-up | Excellent: Consumes accumulated burst up to bucket capacity | May overwhelm downstream dependencies if burst capacity is misconfigured too high. |
| Leaky Bucket | Moderate: Queue holding requests | Moderate: Constant rate dequeue scheduling | Smooths bursts into a stable, deterministic flow | Can drop real-time requests if the buffer fills during traffic spikes. |
| Fixed Window Counter | Extremely Low: Single atomic integer per window | Lowest: Atomic counter increments | Poor: Allows up to 2x configured burst limit at window boundaries | Boundary clustering can overload backend systems at window edges. |
| Sliding Window Log | High: O(M) where M is the number of logged requests | High: Sorting, pruning, and counting timestamps | Accurate: Prevents window boundary traffic spikes | High memory usage; vulnerable to memory exhaustion under traffic surges. |
| Sliding Window Counter | Low: Counters for current and previous windows | Low: Weighted average calculation across windows | Good: Smooths out bursts with low operational overhead | Approximation based on an assumed uniform distribution in the prior window. |
Circuit Breaker State Machine Mechanics
When an upstream dependency fails or degrades, continuous synchronous requests exhaust caller resources. A circuit breaker monitors call success rates and isolates broken services:
[ Failure Threshold Exceeded ]
+------------------------------------------------+
| |
v |
+--------+ Success Count > Threshold +-------------+ Failure Occurs +------+
| CLOSED | <---------------------------- | HALF-OPEN | <--------------- | OPEN |
+--------+ +-------------+ +------+
^ |
| Sleep Window Elapsed |
+------------------------------------------------------------+
- Closed: Normal operating state. Requests pass through directly. Failures increment an internal counter. If the failure rate crosses a configured threshold (e.g. 50% over a 10-second rolling window), the breaker trips into the Open state.
- Open: All incoming requests fail fast immediately, returning a local fallback response (e.g. cached default, empty response, or HTTP 503) without touching the degraded dependency. This relieves network and CPU pressure on downstream systems.
- Half-Open: After a cool-down sleep duration (e.g. 30 seconds), the circuit allows a small, throttled trial percentage of real requests through. If these canary requests succeed, the breaker resets to Closed. If any canary fails, the breaker returns to Open for another cool-down window.
Production Circuit Breaker Implementation
The following Go implementation shows a production-grade circuit breaker managing state transitions, concurrency safety, and recovery timeouts.
package resiliency
import (
"errors"
"sync"
"time"
)
type State int
const (
StateClosed State = iota
StateOpen
StateHalfOpen
)
var ErrCircuitOpen = errors.New("circuit breaker is open; call rejected")
type CircuitBreaker struct {
mu sync.Mutex
state State
failureThreshold int
failureCount int
successThreshold int
successCount int
timeout time.Duration
lastStateChange time.Time
}
func NewCircuitBreaker(threshold int, timeout time.Duration, successThreshold int) *CircuitBreaker {
return &CircuitBreaker{
state: StateClosed,
failureThreshold: threshold,
successThreshold: successThreshold,
timeout: timeout,
lastStateChange: time.Now(),
}
}
func (cb *CircuitBreaker) Execute(action func() error) error {
cb.mu.Lock()
now:= time.Now()
// Check if Open state has timed out, transitioning to Half-Open
if cb.state == StateOpen && now.Sub(cb.lastStateChange) > cb.timeout {
cb.state = StateHalfOpen
cb.failureCount = 0
cb.successCount = 0
cb.lastStateChange = now
}
if cb.state == StateOpen {
cb.mu.Unlock()
return ErrCircuitOpen
}
cb.mu.Unlock()
// Execute target action
err:= action()
cb.mu.Lock()
defer cb.mu.Unlock()
if err!= nil {
cb.failureCount++
if cb.state == StateHalfOpen || cb.failureCount >= cb.failureThreshold {
cb.state = StateOpen
cb.lastStateChange = time.Now()
}
return err
}
// Handle successful execution
if cb.state == StateHalfOpen {
cb.successCount++
if cb.successCount >= cb.successThreshold {
cb.state = StateClosed
cb.failureCount = 0
cb.successCount = 0
cb.lastStateChange = time.Now()
}
}
return nil
}
Frequently Asked Questions
What are the most critical latency numbers to memorize for back-of-the-envelope calculations?
Key baselines include L1 cache references at 1 nanosecond, RAM access at 100 nanoseconds, NVMe SSD reads at 10 to 50 microseconds, intra-datacenter round trips at 500 microseconds, and cross-continent network round trips at 100 to 150 milliseconds. Use these building blocks to estimate realistic end-to-end API latency budgets.
How should an engineer time-box a 45-minute technical architecture round?
Spend 5 minutes clarifying scope and constraints, 10 minutes drafting high-level architecture and API contracts, 20 minutes deep-diving into database scaling and component interactions, and 10 minutes addressing edge cases, failure domains, bottlenecks, and observability instrumentation.
What is the primary difference between CAP theorem and PACELC?
CAP addresses system behavior solely during network partitions, requiring a choice between consistency and availability. PACELC extends CAP by stating that even during normal operation (Else), an engineer must explicitly trade off latency versus consistency when replicating data across distributed nodes.
How do you prevent the thundering herd problem in distributed caches?
Prevent cache stampedes using distributed mutex locking (such as Redlock) to allow only one worker to query storage, configuring probabilistic early recomputation (XFetch), or serving stale cache values in the background while an asynchronous worker regenerates the payload.
Designing resilient distributed architectures requires balancing latency trade-offs, consensus requirements, and storage characteristics. Every architectural decision involves fundamental trade-offs: picking between the immediate consistency of a B+ Tree RDBMS or the write-optimized throughput of an LSM-Tree; deciding between CP and AP within CAP and PACELC constraints; and balancing lazy-loaded cache simplicity against write-behind performance.
Use this reference framework as an operational checklist: anchor capacity planning to hardware physics, structure architectural presentations methodically, and isolate failure domains at the edge. Building scalable systems is not about chasing theoretical perfection, but about matching your operational requirements to the right architectural patterns while proactively engineering for real-world hardware failures.