Interview Setup
Interview Prompt
Design a distributed key-value store modeled after Amazon Dynamo and Apache Cassandra. Support put, get, and delete operations by key, masterless replication across nodes, and configurable durability. Target web-scale workloads with partition tolerance.
Clarifying Questions (ask before designing)
| Question | Why it matters |
|---|---|
| What level of consistency does the client require between eventual consistency and linearizable reads? | Determines read and write quorum sizes, leaderless coordination models, and conflict resolution strategies. |
| What is the workload distribution across reads and writes, and what is the typical value size? | LSM trees optimize write-heavy append workloads, whereas B-trees benefit read-heavy small-record lookups. Value size governs compaction overhead and network replication costs. |
| Do we need compare-and-swap semantics, multi-key transactions, or single-key operations only? | Multi-key transactions require distributed consensus protocols, whereas single-key operations execute locally within partition token ranges. |
| Should consistency levels be tunable per request or enforced as a cluster-wide policy? | Supporting per-request consistency levels such as ONE, QUORUM, and ALL allows systems to satisfy diverse service level agreements across distinct workloads. |
Scope
In scope
- Consistent hashing and virtual node partitioning
- Put, get, and delete operations with quorum reads and writes
- Failure detection and fault recovery across nodes and network partitions
- On-disk LSM persistence with Write-Ahead Logging and SSTable compaction
- Conflict resolution for concurrent writes using vector clocks and Last-Write-Wins
Functional Requirements
Clarify the exact operational scope with the interviewer. Core operations focus on key-based access, alongside expiration policies, versioning, and replication behaviors.
put(key, value): Store a key-value pair, overwriting any existing value for the key.get(key): Retrieve the value associated with a given key, returning null if the key is not present.delete(key): Remove a key-value pair by writing a deletion tombstone marker.- Support keys up to 256 bytes and values up to 10 KB in size.
- Support versioning and conflict resolution using vector clocks or Last-Write-Wins.
- Support configurable time-to-live (TTL) expiration on individual keys.
- Support range scans across ordered token ranges where appropriate.
Non-Functional Requirements
For a Dynamo-style distributed store, evaluations focus on partition tolerance, high write throughput, and tunable consistency rather than relational multi-table transactions.
- High Availability: Operate as an AP system under the CAP Theorem, prioritizing write availability during partial network outages by utilizing sloppy quorums when requested consistency levels allow.
- Tunable Consistency: Allow clients to specify consistency levels per request, ranging from eventual consistency (ONE) to quorum overlap consistency (QUORUM).
- Low Latency: Achieve sub-10-millisecond p99 latency for both read and write operations.
- Horizontal Scalability: Scale storage and throughput linearly by adding commodity nodes to the cluster.
- Durability: Acknowledged writes are designed to survive individual node failures through synchronous local Write-Ahead Logging and quorum replication across distinct failure domains.
- Partition Tolerance: Maintain cluster operation and data availability during inter-rack or cross-zone network partitions.
- Automated Fault Recovery: Detect node outages and synchronize divergent data automatically through gossip and anti-entropy protocols.
Capacity Estimations
Establish capacity figures before finalizing cluster node sizing and storage parameters. Peak query rates determine coordinator throughput requirements, while total data volume dictates the number of storage nodes.
| Metric | Calculation | Value |
|---|---|---|
| Total logical data | Given | 1 PB |
| Average key size | Given | ~100 bytes |
| Average value size | Given | 1 KB |
| Storage per entry | 100B + 1 KB | ~1.1 KB |
| Replication factor 3 | 1 PB x 3 | ~3 PB physical |
| Writes/sec (peak) | Given | 1M |
| Reads/sec (peak) | Given | 3M |
| Nodes (10 TB each) | 3 PB ÷ 10 TB | ~300 nodes |
Storage and Node Sizing Calculations
- Raw Logical Volume: 1 PB of logical data with an average entry size of 1.1 KB represents approximately 909 billion keys.
- Replicated Physical Footprint: With a replication factor of N=3, raw storage requirements reach 3 PB before accounting for compaction overhead.
- Node Allocation: Assuming commodity storage nodes equipped with 10 TB NVMe SSDs and maintaining a 70% storage utilization ceiling to accommodate compaction spikes, the cluster requires approximately 300 to 400 storage nodes.
Architecture Diagram
Interview strategy: State your CAP Theorem alignment upfront. Emphasize that the system implements a masterless AP architecture with tunable quorums rather than centralized single-leader transactions.
Client requests arrive at any cluster node, which acts as the coordinator for that operation. The coordinator hashes the requested key onto the 128-bit consistent hash ring, identifies the primary owner node and its next N-1 successor replica nodes, and orchestrates the quorum interaction.
On the write path, data is appended to the local Write-Ahead Log for crash recovery and inserted into the in-memory MemTable before being dispatched to replica nodes. On the read path, the coordinator queries R replica nodes in parallel, evaluates versions using vector clocks or timestamps, returns the winning value to the client, and initiates background read repair for any stale replicas.
Component Deep Dives
The masterless key-value architecture combines consistent hashing, an LSM storage engine, gossip membership, and quorum protocols to balance throughput, availability, and durability.
Consistent Hashing with Virtual Nodes
Consistent hashing maps keys and physical nodes onto a shared 128-bit ring, ensuring that cluster resizing operations only migrate a minimal fraction of token ranges.
- Ring Layout: The keyspace spans from 0 to 2128 - 1 using cryptographic hashing (such as MD5 or MurmurHash3) applied to the partition key.
- Virtual Nodes (vnodes): Each physical node is assigned multiple virtual tokens (typically 128 to 256 vnodes) spread across the ring. This balances token distribution, handles heterogeneous hardware capacities, and ensures smooth load redistribution when nodes join or leave.
- Replica Selection: A key maps to a point on the ring, and the first physical node encountered clockwise serves as the primary replica. The next N-1 distinct physical nodes clockwise form the replica set.
- Rack Awareness: The ring topology maps replica assignments across distinct server racks and availability zones to prevent concurrent outages from single-switch or power failures.
Storage Engine: LSM-Tree Architecture
Log-Structured Merge-trees optimize write performance by converting random updates into sequential disk operations across memory and disk structures.
Write Workflow
- Write-Ahead Log (WAL): The mutation is appended to an on-disk sequential log and flushed to provide durable local crash recovery.
- Entry format:
CRC (4B) | Timestamp (8B) | Key Length (4B) | Value Length (4B) | Key | Value - Recovery: On unexpected node restart, the system replays WAL records to reconstruct the active MemTable state locally. If a node suffers unrecoverable local disk loss, the cluster reconstructs its partitions by streaming data from surviving replica peers.
- Lifecycle: WAL segments are truncated once their corresponding MemTable has fully flushed to disk.
- Entry format:
- In-Memory MemTable: The record is inserted into a concurrent skip list that maintains keys in sorted order in memory with O(log n) time complexity.
- MemTable Flushing: When the active MemTable reaches its memory threshold (such as 64 MB), it transitions to an immutable state and a new active MemTable is initialized. A background thread writes the immutable MemTable to disk as a sorted, immutable SSTable file.
Read Workflow
- Check the active in-memory MemTable for the most recent version of the key.
- Check any pending immutable MemTables awaiting disk flush.
- Evaluate the Bloom filter for each relevant on-disk SSTable. If the filter indicates the key is definitely not present, skip disk access for that file entirely.
- If the Bloom filter indicates the key may be present, perform a binary search on the SSTable sparse index to locate the target data block offset on disk.
- Read and decompress the data block from disk, scan for the exact key match, and return the record with the latest timestamp or vector clock.
SSTable Internal Structure
- Data Blocks: Compressed sequential blocks containing sorted key-value pairs along with mutation timestamps and deletion markers.
- Sparse Index Block: Contains offsets for every Nth key (such as every 4 KB boundary), enabling binary searches across data blocks with minimal memory overhead.
- Filter Block: Contains a Bloom filter configured with 10 bits per key, providing an illustrative 1% false-positive rate. This significantly reduces unnecessary disk seeks for absent keys, with the exact reduction depending on key access patterns and the number of SSTables consulted.
- Footer Block: A fixed-size 48-byte trailer containing byte offsets and lengths for the index and filter blocks, along with format versions and magic verification bytes.
Compaction Strategies
- Size-Tiered Compaction Strategy (STCS): Merges SSTables of similar file sizes into larger files when the table count exceeds a threshold. This approach minimizes write amplification and works well for write-heavy workloads, though it requires up to 50% temporary free disk space.
- Leveled Compaction Strategy (LCS): Organizes SSTables into fixed levels where each level is 10 times the size of the previous level, and key ranges within levels above Level 0 are non-overlapping. This approach lowers read amplification and space overhead at the expense of higher write amplification during background merges.
Replication and Quorum Models
Replication ensures fault tolerance, while tunable quorums allow applications to balance consistency and latency according to their operational needs.
- Replication Factor (N): Specifies the total number of distinct physical nodes that store copies of each key-value pair, with N=3 being a common production choice. The exact factor is chosen based on durability targets, availability requirements, storage costs, and failure domain boundaries.
- Quorum Overlap (R + W > N): Choosing R and W such that R + W > N ensures that the read quorum and write quorum intersect in at least one replica. This significantly improves read-after-write visibility under the protocol assumptions. However, quorum overlap alone does not guarantee universal linearizability, as concurrent writes, conflict resolution, clock skew, and failure timing still govern final data convergence.
- Tunable Quorum Modes:
W=1, R=1: Delivers lowest latency and highest availability with eventual consistency.W=2, R=2 (with N=3): Provides quorum overlap across standard read-write workloads.W=3, R=1: Provides fast read access with synchronous write durability across all replicas.LOCAL_QUORUM: Restricts quorum evaluation to the local datacenter to avoid cross-region network latency.
Gossip Protocol and Failure Detection
A decentralized peer-to-peer gossip protocol maintains cluster membership state and detects node outages without requiring a centralized master.
- State Exchange: Every second, each node selects a random peer and exchanges heartbeat sequences, token ownership maps, and schema metadata.
- Convergence: Information propagates across a cluster of N nodes in O(log N) gossip rounds, achieving global cluster state synchronization within seconds.
- Phi Accrual Failure Detector: Rather than using binary up or down timeouts, nodes calculate a continuous suspicion value ($\Phi$) based on historical heartbeat arrival distributions. A node is declared unreachable when $\Phi$ crosses a configured threshold (such as $\Phi \ge 8$), adapting dynamically to network latency fluctuations.
Hinted Handoff for Temporary Outages
Hinted handoff allows write operations to succeed even when one of the designated replica nodes is temporarily unreachable, preserving write availability in AP mode.
- Mechanism: If a coordinator fails to contact a designated replica during a write, it writes the payload to a healthy substitute node along with a hint containing the target node ID, key, value, and timestamp.
- Replay Delivery: The substitute node monitors gossip heartbeats for the target replica. Once the primary replica returns to service, the substitute streams the buffered hinted writes and deletes the hints upon verified acknowledgment.
- Hint Expiration Window: Hints are retained for a bounded duration (typically 3 hours). If the target node remains down past the TTL window, the hints are discarded and the node is reconciled via Merkle tree anti-entropy repair.
Anti-Entropy Repair Using Merkle Trees
Anti-entropy repair runs as a background process to detect and reconcile data divergence across replicas caused by dropped hints or extended network partitions.
- Tree Structure: Each node maintains a binary Merkle tree for each token range it manages. Leaf nodes represent hashes of individual key-value pairs, while parent nodes represent cryptographic hashes of their children.
- Comparison Workflow: Two replicas compare root hashes for a shared token range. If root hashes match, the replicas are identical and no data transfer occurs. If root hashes differ, nodes traverse down tree branches to pinpoint the exact divergent subtrees.
- Bandwidth Optimization: For a token range with 1 million keys, comparing Merkle tree branches identifies divergent keys in O(log N x diff_count) operations, synchronizing only out-of-date records without transferring the full dataset.
Conflict Resolution Mechanisms
Concurrent writes across distributed replicas require deterministic conflict resolution strategies to resolve divergent versions.
- Vector Clocks (Amazon Dynamo):
- Each mutation attaches a vector clock tracking logical version counters per node (such as
node_A: 3, node_B: 1). - If one version causally dominates the other, the dominant version is selected automatically. If versions are concurrent, the system preserves both siblings and returns them to the client application for domain-specific reconciliation.
- To prevent unbounded vector clock growth, clocks are pruned using threshold timestamps when the entry count exceeds configured limits.
- Each mutation attaches a vector clock tracking logical version counters per node (such as
- Last-Write-Wins (Apache Cassandra):
- Each write attaches a timestamp, and the mutation with the highest timestamp is chosen as the authoritative winner while older versions are discarded.
- To prevent physical clock skew from silently overwriting newer data, production implementations use Hybrid Logical Clocks (HLC) that combine physical wall-clock time with logical sequence counters.
Coordinator Node Workflow
Any node in the cluster can accept a client request and coordinate its execution across the appropriate replica set.
- Write Coordination:
- The coordinator appends the mutation to its local Write-Ahead Log and inserts it into its active MemTable.
- The coordinator forwards the write request concurrently to the remaining N-1 replica nodes responsible for the key's token range.
- The coordinator waits for W total acknowledgments (including its own local write).
- Once W acknowledgments are received, the coordinator returns success to the client. If a replica is unreachable, a hinted handoff payload is stored on an available node.
- Read Coordination and Read Repair:
- The coordinator sends read requests in parallel to R replica nodes.
- The coordinator collects the responses and compares their version metadata.
- The latest value is returned immediately to the client to satisfy the read latency SLA.
- If any replica returned an older version, the coordinator asynchronously dispatches the latest record to the stale replicas in the background via Read Repair.
API Design
The store exposes a minimal, high-throughput key-value interface alongside internal node-to-node replication and anti-entropy repair contracts.
Client API Interface
External client contract supporting version-aware, tunable-consistency operations:
export interface RequestContext {
vectorClock?: Record<string, number>;
timestamp?: number;
consistencyLevel: "ONE" | "QUORUM" | "ALL" | "LOCAL_QUORUM";
}
export interface PutRequest {
key: string;
value: string;
context: RequestContext;
}
export interface GetResponse {
value: string | null;
context: RequestContext;
siblings?: string[];
}
export interface KeyValueClientApi {
put(request: PutRequest): Promise<{ success: boolean }>;
get(key: string, consistency: RequestContext["consistencyLevel"]): Promise<GetResponse>;
delete(key: string, context: RequestContext): Promise<{ success: boolean }>;
}Internal Replication and Cluster APIs
Inter-node communication protocols for data replication, hinted handoff, and background anti-entropy repairs:
export interface KeyValueInternalApi {
replicate(
key: string,
value: string,
timestamp: number,
vectorClock: Record<string, number>
): Promise<{ ack: boolean }>;
hintedHandoff(
targetNodeId: string,
key: string,
value: string,
timestamp: number
): Promise<{ ack: boolean }>;
antiEntropyRepair(
tokenRange: string,
merkleRootHash: string
): Promise<{ syncedKeys: number }>;
}Common Error Responses
400 Bad Request: invalid input, missing required fields, or malformed JSON payload 401 Unauthorized: missing or invalid authentication token or API key 403 Forbidden: authenticated caller lacks required permissions for this resource 404 Not Found: requested resource ID does not exist 409 Conflict: duplicate write or version conflict, retry with a unique idempotency key 422 Unprocessable Entity: syntactically valid request failed semantic business validation 429 Too Many Requests: rate limit quota exceeded, client should honor Retry-After header 500 Internal Error: unexpected server failure, retry safely with an idempotency key 503 Service Unavailable: downstream dependency is unavailable or overloaded, retry with exponential backoff
Data Model
The physical storage model centers on immutable on-disk SSTable files, in-memory MemTables, Write-Ahead Logs, and version metadata.
On-Disk SSTable Layout
SSTable files consist of contiguous, independently compressed blocks with sparse indexing:
- Data Blocks: Sequential key-value pairs sorted lexicographically by partition key. Each entry contains key bytes, mutation timestamp, value length, and payload bytes.
- Index Block (Sparse): Maps the first key of every data block to its physical byte offset within the SSTable file, enabling binary search lookups without loading the entire index into memory.
- Filter Block: Bloom filter bit arrays corresponding to data blocks, checked before issuing disk seeks.
- Meta Block: Stores table compression codecs, dictionary tables, and creation metadata.
- Footer: A fixed 48-byte structure at the end of the file containing exact offsets and lengths for the index and meta blocks, along with magic verification bytes.
In-Memory MemTable Structure
Data Structure: Skip List (O(log n) concurrent inserts, lookups, and range scans) Capacity Threshold: 64 MB Lifecycle: When full, becomes an immutable MemTable and a new active MemTable is initialized while a background worker flushes the frozen instance to an on-disk SSTable.
Write-Ahead Log Entry Format
| CRC (4B) | Timestamp (8B) | Key Length (4B) | Value Length (4B) | Key (variable) | Value (variable) |
Vector Clock Payload Structure
{
"key": "user:123:profile",
"value": "{\"name\": \"Alice\"}",
"vector_clock": {
"node_A": 3,
"node_B": 1,
"node_C": 2
},
"timestamp": 1710320000
}Fault Tolerance
Masterless architectures achieve resilience through decentralized redundancy, automated failover routing, and eventual convergence mechanisms.
Resilience Mechanisms Summary
| Mechanism | Operational Function |
|---|---|
| Replication Factor (N=3) | Stores data across three distinct physical nodes to tolerate up to two simultaneous node outages for eventual reads. |
| Write-Ahead Logging (WAL) | Appends mutations to a local sequential disk log with fsync before acknowledging writes, ensuring local crash recovery before MemTable flushes. |
| Hinted Handoff | Buffers writes on substitute nodes when designated replicas are unreachable, replaying data once nodes recover. |
| Anti-Entropy (Merkle Trees) | Compares cryptographic hash trees between replicas in the background to detect and synchronize divergent key ranges. |
| Read Repair | Reconciles divergent version data in the background whenever a quorum read detects stale replica responses. |
Failure Scenarios and Mitigations
1. Transient Node Outage
- The coordinator buffers mutations intended for the offline replica as hinted handoffs on surviving nodes.
- Gossip failure detectors broadcast node unavailability across the cluster.
- Read requests continue to succeed from remaining replicas as long as active replicas satisfy read quorum requirements.
2. Permanent Node Loss and Hardware Decommissioning
- The node is marked dead after exceeding the gossip suspicion threshold timeout.
- Surviving nodes redistribute the dead node's virtual token ranges across remaining cluster members.
- New replicas stream missing token ranges from surviving replica peers to restore the replication factor of N=3.
3. Network Partitions
- Under AP configuration, nodes on both sides of the partition accept reads and writes using sloppy quorums.
- When the network partition heals, background Merkle tree anti-entropy repairs and vector clock reconciliations merge divergent histories.
4. SSTable Block Corruption
- CRC32 checksums validate every SSTable data block and WAL entry during disk reads.
- If corruption is detected, the coordinator rejects the corrupted block and streams a verified replacement block from healthy replica peers.
5. Compaction Backpressure and I/O Starvation
- Excessive flush rates can lead to compaction backlogs and high read amplification.
- Mitigations include rate-limiting compaction I/O bandwidth, prioritizing client reads and writes, and dynamically increasing worker thread pools.
Additional Considerations
Advanced production considerations cover probabilistic data structures, tombstone lifecycles, and hot key mitigation strategies.
Bloom Filter Optimization
- Each SSTable includes an in-memory Bloom filter to prevent unnecessary disk seeks for keys not present in that file.
- Configured with 10 bits per key and 7 hash functions, yielding an approximate 1% false-positive rate.
- Because Bloom filters have zero false negatives, a negative result guarantees the key is absent, allowing queries to safely skip checking that SSTable file.
Tombstones and Garbage Collection Lifecycles
- Delete operations append a deletion tombstone record containing a deletion timestamp rather than immediately removing data in place.
- Tombstones persist across a configurable grace period (typically 10 days, corresponding to
gc_grace_seconds) to ensure all replicas and offline nodes learn of the deletion. - Once the grace period expires, subsequent compaction cycles permanently purge the tombstone and all older shadowed versions from disk.
Hot Key Mitigation
- Excessive traffic to a single popular key can saturate the responsible replica nodes.
- Mitigations include placing an in-memory cache layer in front of the cluster, distributing reads across all N replicas, and prefixing partition keys with shard identifiers for high-write hot keys.
Tunable Consistency Profiles
| Configuration | Consistency Model | Typical Use Case |
|---|---|---|
| W=1, R=1 | Fastest response with eventual consistency | Ephemeral session state and click tracking |
| W=2, R=2, N=3 | Quorum overlap (R + W > N), ensuring read and write quorums intersect in at least one replica | User profiles, shopping carts, and inventory counts |
| W=3, R=1 | Fast low-latency reads with durable synchronous writes | Write-rarely, read-frequently reference data |
| W=1, R=3 | Fast writes with heavyweight quorum reads | High-volume logging and audit trails requiring verified reads |
Related Problems and Concepts
Storage engine internals and partition management connect directly to patterns in Pastebin and Design a URL Shortener. Review foundational concepts in Consistent Hashing, SQL vs NoSQL, CAP Theorem and PACELC, Replication, Failover, and Leader Election, Sharding and Partitioning, and WAL and Data Durability.
Interview Walkthrough
- 25-Minute Interview Strategy
Focus on the core AP architecture, consistent hashing, and LSM storage engine mechanics before diving into operational edge cases.
- CAP Theorem and AP System Choice (3 min)
- Consistent Hashing with Virtual Nodes (5 min)
- Replication Quorums and Tunable Consistency (5 min)
- LSM Storage Engine Write and Read Paths (7 min)
- Hinted Handoff and Fault Tolerance (5 min)
- Frame the problem as a write-heavy, partition-tolerant distributed store, clarifying why web-scale key-value workloads generally choose availability over linearizable transactions.
- Explain Consistent Hashing and virtual nodes upfront to demonstrate how data distribution remains balanced during cluster resizing.
- Contrast LSM-trees with B-trees using Back-of-the-Envelope Estimation on write amplification, emphasizing why sequential disk appends excel at millions of writes per second.
- Detail the quorum formula (R + W > N) and explain how sloppy quorums with hinted handoffs maintain write availability during node outages.
- Discuss conflict resolution options by contrasting vector clocks with Last-Write-Wins and Hybrid Logical Clocks.
- Explain how Bloom filters and in-memory MemTables optimize the read path to keep latency predictable under Caching Patterns and Invalidation.
Engineering Trade-offs
Architectural decisions in distributed storage balance write throughput, read latency, consistency guarantees, and operational complexity.
Storage Engine: LSM-Tree vs B-Tree
Selecting the on-disk storage engine represents a fundamental trade-off between write amplification and point lookup latency. The complexity labels in the table below provide simplified intuition for interview comparison rather than exact production costs. Real-world behavior depends on indexing depth, Bloom filter efficiency, caching, level structures, compaction schedules, and storage engine implementations.
| Evaluation Factor | LSM-Tree (This Architecture) | B-Tree (InnoDB / PostgreSQL) |
|---|---|---|
| Write performance | O(1) amortized via sequential disk append | O(log N) per write with random page I/O |
| Read performance | O(N) worst case across multiple SSTables | O(log N) predictable single tree traversal |
| Write amplification | High (10x to 50x rewritten during compaction) | Low to moderate (in-place page updates) |
| Read amplification | Higher (may check multiple SSTable levels) | Low (single tree traversal per lookup) |
| Space amplification | Moderate (temporary versions during compaction) | Low (in-place updates with page fragmentation) |
| Disk I/O pattern | Sequential append-only writes | Random I/O page updates |
Why LSM-Tree is selected: Distributed key-value workloads are predominantly write-heavy, ingesting continuous mutations and session updates. The sequential append architecture of LSM-trees provides substantially higher write throughput than in-place B-tree updates. Read latency penalties are mitigated using Bloom filters to filter out the vast majority of absent key lookups and serving recent data directly from the in-memory MemTable.
When B-Tree is preferred: Read-heavy relational workloads requiring predictable single-digit point reads and complex range scans where write volume is moderate.
CAP Theorem: AP vs CP Architecture
Choosing between availability and consistency defines system behavior during network partitions.
| Dimension | AP Configuration (Dynamo / Cassandra) | CP Configuration (etcd / ZooKeeper) |
|---|---|---|
| Partition behavior | Both partition halves continue accepting reads and writes | Minority partition rejects writes to prevent split-brain |
| Consistency guarantee | Eventual consistency with asynchronous reconciliation | Linearizable strong consistency via Raft or Paxos |
| Availability SLA | High write availability via sloppy quorum (subject to requested consistency level) | Degraded availability during leader election or partitions |
| Best use cases | Shopping carts, user sessions, telemetry, and profiles | Distributed locks, configuration catalogs, and leader election |
Why AP is chosen for general key-value storage: Most web-scale services (such as shopping carts and session registries) prefer accepting writes with temporary staleness over returning errors to users during network disruptions.
Conflict Resolution: Vector Clocks vs Last-Write-Wins
Resolving divergent concurrent writes requires balancing data safety against client implementation complexity.
| Approach | Strengths | Weaknesses | Mitigation Strategy |
|---|---|---|---|
| Vector Clocks (Amazon Dynamo) |
|
| Prune oldest vector clock entries when size exceeds a threshold (such as 10 entries) |
| Last-Write-Wins (Apache Cassandra) | Simple evaluation with fixed-size timestamp metadata and no client merge burden | Can silently discard concurrent writes if timestamps collide or drift | Use Hybrid Logical Clocks (HLC) combining physical time with monotonic counters to avoid clock skew errors |
Recommendation: Vector clocks demonstrate architectural depth in interviews. In enterprise production environments, Last-Write-Wins paired with Hybrid Logical Clocks is widely favored due to lower operational overhead.
Virtual Nodes: Uniformity vs Management Overhead
Virtual nodes solve ring token imbalance and simplify cluster rebalancing operations.
| Architecture | Key Distribution | Rebalancing Behavior | Hardware Heterogeneity |
|---|---|---|---|
| Single Token per Physical Node | High variance in token ownership, leading to hotspot nodes | Removing a node dumps all its data onto a single neighbor node | Cannot easily adapt to servers with differing CPU and storage capacities |
| Virtual Nodes (128-256 vnodes/node) ⭐ | Near-uniform statistical distribution across all nodes via the law of large numbers | Removing a node redistributes its vnodes evenly across all surviving cluster members | Assign proportional vnode counts to hardware tiers (such as 256 vnodes for large nodes and 64 for small nodes) |
Compaction Strategy: Size-Tiered vs Leveled
Compaction manages disk space and read performance as SSTables accumulate on disk.
| Compaction Strategy | Write Amplification | Read Amplification | Space Overhead | Recommended Workload |
|---|---|---|---|---|
| Size-Tiered Compaction (STCS) | Low (merges similarly sized files) | High (many overlapping SSTables to check) | High (up to 50% temporary disk headroom required) | Write-heavy workloads where disk capacity is ample |
| Leveled Compaction (LCS) | High (rewrites entire levels during promotion) | Low (non-overlapping key ranges within levels) | Low (approximately 10% disk overhead) | Read-heavy workloads requiring predictable read latency |
Recommendation: Start with Size-Tiered Compaction for high-throughput write ingestion. Transition to Leveled Compaction if read p99 latency degrades due to scanning excessive SSTables.
Quorum Execution: Strict vs Sloppy Quorum
Sloppy quorums provide high availability during temporary node failures at the cost of short-term read consistency.
| Model | Execution Rule | Availability Guarantee | Consistency Guarantee |
|---|---|---|---|
| Strict Quorum | Writes must be acknowledged exclusively by designated primary replicas for the token range | Fails writes if fewer than W primary replicas are available | Ensures strict replica membership without temporary misplaced writes |
| Sloppy Quorum (Hinted Handoff) ⭐ | Writes are accepted by any W available healthy nodes, buffering hints for offline replicas | Maintains high write availability during multi-node outages when configured consistency thresholds are met | Eventual consistency until hinted handoffs and anti-entropy repairs complete |
Review
How helpful was this walkthrough?
Click a star to rate. We actively use this feedback to refine and update our system design content.
Discussion
Share your thoughts, ask questions, or help others.