Design a Distributed Key-Value Store
1. Problem Statement & Scope
System Mission
Design a distributed, highly available, fault-tolerant, horizontally partitioned Key-Value Store (combining the architectural principles of the Amazon Dynamo paper and Apache Cassandra) that stores tens of terabytes of structured and semi-structured key-value data today, serves 250,000 requests per second at peak with single-digit millisecond latency, scales linearly by adding nodes, and survives the loss of a whole Availability Zone without an outage.
Functional Requirements
put(key, value, options): Store or mutate an opaque binary or JSON payload (up to ) indexed by a string partition key. Supports optional TTL, idempotency tokens, and conditional expression evaluation (e.g.,attribute_not_exists(key)orversion = 4).get(key, options): Retrieve the value and version metadata associated with a partition key, at a tunable read consistency level:ONE(eventual),QUORUM, orLOCAL_QUORUM(2 of the 3 replicas in the caller's Region; the same asQUORUMuntil a second Region exists, Section 8.6).delete(key): Soft-delete a key-value record by writing a tombstone record with a monotonic timestamp, later reclaimed asynchronously by background compaction.batch_get(keys)/batch_write(items): Execute high-throughput multi-key read and write operations across different partitions with partial success handling (up to 100 keys per read batch, 25 items or 16 MB per write batch).- Tunable Consistency Configuration: Let clients select consistency levels per request or per table: so every read overlaps the latest acknowledged write, or for the lowest-latency eventual consistency.
Architectural Note — Why Is There No POST(key, value) in Functional Requirements?
In distributed key-value storage systems (such as Amazon DynamoDB PutItem, Redis SET, Amazon S3 PutObject, and Apache Cassandra INSERT), single-record creation and mutation are unified under put(key, value, options) using idempotent upsert semantics ():
- Client-Supplied Partition Keys: Unlike relational CRUD APIs where the server generates auto-incrementing IDs or UUIDs via
POST /resources, a key-value store's primary access path is directly addressed by the client-provided key (/v1/kv/{key}). - Idempotency & Safe Network Retries: Over unreliable distributed networks, client timeouts frequently occur after a write has committed to storage replicas. Because
PUTis strictly idempotent, client SDKs can safely retry timed-out requests with exponential backoff without risking duplicated records or phantom state corruption. - Preconditions Provide "Create-Only" Semantics: If a caller requires create-only behavior (i.e. aborting if the key already exists), this is achieved through conditional expressions (e.g.,
condition_expression: "attribute_not_exists(key)"or HTTPIf-None-Match: *). If the key exists, the storage engine rejects the write with412 Precondition Failed, completely eliminating the need for a separatePOSTverb. (Note:POSTis reserved at the HTTP API layer exclusively for composite batch actions—POST /v1/kv:batch-getandPOST /v1/kv:batch-write—to overcome URI length limits onGETand model partial batch processing across multiple partitions).
Value Data Shapes: Structured vs. Semi-Structured
The store enforces no schema on values: each one is an opaque binary or JSON payload of up to , so it holds both structured data and semi-structured data. What differs is where the schema is enforced and what that costs.
1. Structured Key-Value Data
In a structured approach, the values stored alongside the key strictly conform to a predefined schema with fixed fields, types, and constraints.
- Schema Rigidity: Every record must adhere to a strict definition. Relational tables and Apache Cassandra's CQL tables (typed columns) enforce it inside the database. Amazon DynamoDB does not: it types only the key attributes, and every other attribute is free-form, which makes it semi-structured. This store enforces no value schema, so the schema lives in the client, as a fixed serialization format such as Protocol Buffers or Avro that the writer validates.
- Data Integrity & Efficiency: Fixed fields allow optimized binary serialization, predictable memory layouts, and efficient field-level retrieval. Updating one field without rewriting the entire payload needs an engine that stores each field separately, as Cassandra does with per-column cells. Here
putreplaces the whole value, so a one-field change rewrites the full value and its version. - Common Use Cases: Financial ledgers, user profile records, transaction logs, and high-frequency time-series metrics.
textKey: "user:1001" Value Schema: { id: INT64, email: STRING, status: ENUM, updated_at: TIMESTAMP }
2. Semi-Structured Key-Value Data
Semi-structured data introduces flexibility by embedding structural metadata (like tags, key-value pairs, or hierarchy) directly inside the stored payload, without requiring a fixed database-wide schema.
- Flexible/Dynamic Schema: Different records under the same dataset can contain entirely different attributes or nested structures. Schema enforcement is handled by the application code rather than the database engine.
- Self-Describing Payloads: Each value carries its own field names (JSON, BSON, MessagePack). That costs more bytes per record than a structured binary encoding, and readers must tolerate missing or unknown fields (schema-on-read).
- Common Use Cases: Product catalogs where each category has different attributes, user preferences and feature flags, session state, and event payloads from many producers.
textKey: "product:88213" Value: { "title": "Running Shoe", "price": 89.99, "attributes": { "color": "red", "size": "10" } } Key: "product:88214" Value: { "title": "Espresso Machine", "price": 349.00, "attributes": { "wattage": 1350 }, "warranty_years": 2 }
Non-Functional Requirements (SLAs/SLOs)
- High Availability:
- per Region across 3 AZs (less than 52.6 minutes of downtime per year).
- for tables replicated to a second Region (less than 5.26 minutes of downtime per year; Section 8.6).
- Latency SLOs:
- Reads (measured at the request router): P95 , P99 .
- Writes (local WAL commit on replicas): P95 , P99 .
- Durability Guarantee: Two layers.
- An acknowledged write is on the fsynced Write-Ahead Log (WAL) of at least replicas in two different AZs, and the third replica catches up through hinted handoff and anti-entropy.
- Continuous backup of SSTables and WAL segments to Amazon S3 (designed for , 11 nines, object durability) bounds data loss to about a minute even if every live replica is lost (Section 10.3).
- The 11-nines figure belongs to S3, not to the 3-replica cluster by itself.
- CAP / PACELC Classification: PA/EL.
- Under a network partition the system favors Availability: the majority side keeps serving
QUORUMrequests, and the minority side still servesONE-level requests but returns503forQUORUM. - In normal operation (Else) it favors Latency over Consistency, unless the caller asks for quorum reads and writes.
- Quorum overlap is not linearizability (Section 3); true compare-and-set uses the conditional-write path (Section 5.4).
- Under a network partition the system favors Availability: the majority side keeps serving
- Scalability: Linear horizontal scaling through consistent-hashing ring partitioning with virtual nodes and no master node on the data path (Section 8.8).
- Operability: Any single node, or a whole AZ, can be replaced or upgraded without downtime or data loss (Sections 8.1 and 10.2).
Out of Scope
Secondary indexes, range scans across keys, and multi-key ACID transactions. A multi-key transaction API would need two-phase commit across partitions (see Two-Phase Commit & Saga); we state it as a follow-up, not a v1 feature.
2. Capacity & Scale Estimation
Traffic Calculations
- Active User Base: 1 Billion Monthly Active Users (MAU), 200 Million Daily Active Users (DAU).
- Total Keys Stored: 10 Billion keys () at launch, growing to 20 Billion over 3 years.
- Query Volume & Read/Write Ratio:
- The workload is read-heavy with an read-to-write ratio.
- Daily request volume (200M DAU, each user's sessions generating key-value operations per day via profile reads, session lookups, and feature-flag checks):
- Average Total QPS:
- Average Read QPS ():
- Average Write QPS ():
- Peak Traffic Factor ( multiplier for daily peaks and marketing campaigns):
Storage Calculations (3-Year Horizon)
- Record Size Breakdown:
- Key size: average.
- Value payload: () average.
- Metadata overhead (version, clock, TTL, flags): .
- Average Record Footprint: .
- Initial Raw Data Volume (10 Billion Keys):
- 3-Year Growth Volume (20 Billion Keys):
- Replication Factor () Across 3 Availability Zones:
- Compaction & Headroom Overhead: LSM-tree leveled compaction needs up to extra space on top of the data for temporary file rewrites and tombstone purge cycles (disks stay under about two-thirds full):
- Compression: ZSTD block compression typically shrinks JSON values 2-4×. We size on uncompressed bytes so compression is pure headroom, never a dependency.
Network Bandwidth
- Ingress Throughput (Write Path):
- Egress Throughput (Read Path):
- Cross-AZ Replication Bandwidth (, writes copied to 2 peer AZs): This traffic is billed per GB; Section 7.3 prices it.
Memory & Cache Sizing (80/20 Pareto Working Set)
- Active Daily Working Set: By the 80/20 rule, of keys account for of read requests: Caching all 4 billion () is not affordable, so we cache the hottest , which serves most of the working set's reads: Spread over the 24 nodes of one AZ (each replica set caches its own hot keys), that is about of block cache per node.
- Bloom Filter Memory Sizing: To avoid costly disk reads on SSTables for keys that do not exist, every SSTable keeps an in-memory Bloom filter with false positive rate (): Across 3 replicas and 72 nodes, that is about per node.
Fleet Sizing & IOPS Provisioning
- Storage Node Sizing:
- Target of EBS gp3 per storage node.
- Required storage nodes across 3 AZs:
- Instance type:
r7g.2xlarge(Graviton3, 8 vCPU, 64 GiB). The 64 GiB holds the 19 GB block cache, about 1 GB of Bloom filters and the MemTables, with room for the OS page cache. - Per-node write IOPS demand during peak: This is pessimistic: WAL group commit merges many writes into one disk flush.
- Per-node read IOPS demand during peak (a quorum read touches 2 replicas, and the digest replica also reads the value to hash it; 90% of reads are served from the MemTable, block cache or OS page cache):
- Storage Specification: Provision EBS
gp3volumes at 6,000 IOPS and 250 MB/s per node: Compaction's large sequential writes are the real extra load, and 250 MB/s covers them within the instance's baseline EBS bandwidth (throughput provisioned above what the instance can push is paid for and never used). Section 8.1 shows why this margin must also cover losing an AZ.
3. AWS-First High-Level Architecture
Synthesizing vector architecture diagram...
Follow a PUT from the top. Route 53 and the NLB spread clients across the "Stateless Request Router Tier" in three AZs; a router keeps no data, so any router can take any request and the tier scales like any web service. In the "Cluster Membership & Hash Ring" panel, routers learn which nodes own which key ranges by gossip, and a phi-accrual detector marks nodes that stop answering. The router hashes the key with Murmur3 and sends the write to the three nodes that own that range, one in each storage panel (for example A1, B1 and C1), so losing a whole AZ still leaves two copies. In the "Storage Engine" panel, each node appends the write to its WAL, adds it to the MemTable, and later flushes it to an SSTable guarded by a Bloom filter; after a crash it replays the WAL. The "Backup & Observability" panel copies every new SSTable to S3 and ships metrics to CloudWatch. The insight: durability comes in two layers, three live replicas for failures you survive online, and S3 for failures you recover from.
Data Flow Walkthrough
- Ingress & Resolution: The client SDK sends the request to Route 53, which resolves to the multi-AZ NLB, which forwards the TCP connection to a stateless Request Router. The router hashes the key with MurmurHash3 (128-bit) and finds the key's preference list: the nodes that own that range on the consistent hash ring, one per AZ.
- Quorum Coordination: The router acts as coordinator. Writes go to all replicas in parallel and succeed after acknowledgments; reads send one full read and one digest read and need answers, repairing any stale replica asynchronously (read repair).
- Local Node Persistence: Each replica appends to its WAL on EBS gp3 (fsynced with group commit), inserts into its concurrent skip-list MemTable, and acknowledges. MemTables flush asynchronously to immutable Level-0 SSTables at .
- Execution Matrix: For every operation's coordinator action, I/O path, guardrails, latency budget and failure recovery, see Section 6.5: Distributed Request & Consistency Execution Matrix.
Tunable Quorum Invariant (): When the sum of write replicas () and read replicas () strictly exceeds the total replication factor (), the Pigeonhole Principle guarantees that at least one replica in any read quorum overlaps with the most recent acknowledged write quorum, so every quorum read sees the latest acknowledged write. When latency is prioritized, yields maximum throughput with bounded eventual consistency.
Where the invariant stops: quorum overlap is not linearizability. Two concurrent writers can still race, and a write that failed after reaching only one replica may or may not appear later. Sloppy quorum breaks the overlap entirely (Section 9.5). Operations that need a true compare-and-set use the conditional-write path in Section 5.4.
Concrete Step-by-Step Request Walkthrough: Tracing Quorum Write & Read Repair
| Step # | Event / Action | Component State | Distributed Transition | Output / Response |
|---|---|---|---|---|
| 1 | Client SDK issues PUT /v1/kv/user_102Consistency=QUORUM (W=2, N=3) | Route 53 resolves to the multi-AZ NLB; TCP is forwarded to a Request Router | Router computes MurmurHash3("user_102") to find the preference list [NodeA1, NodeB1, NodeC1] | Request dispatched in parallel |
| 2 | Router sends the write to all replicas | Router marks the request AWAITING_QUORUM(required=2) | Each replica appends to its EBS gp3 WAL (group-commit fsync) and inserts into its skip-list MemTable | Nodes A1 & B1 return ACK in ; Node C1 times out () |
| 3 | Router receives 2 of 3 ACKs, satisfying | Quorum latch released; a hint for Node C1 is saved on a healthy storage node in the router's AZ | The hint (does not count toward ) is replayed when gossip reports C1 alive again (hint TTL ) | 200 OK(P99 under ) returned to client |
| 4 | Client issues GET /v1/kv/user_102Consistency=QUORUM (R=2, N=3) | Router starts a read: 1 full read (its own AZ) + 1 digest read | Rep1 returns the full value with HLC t=105;Rep2 returns a SHA-256 digest with the older HLC t=98 | Router detects divergence: Rep1 is newer than Rep2 |
| 5 | Router reconciles and repairs | Rep1's version wins; Rep1's value is authoritative | Router answers the client immediately, then sends a non-blocking read repair to Rep2 | 200 OK (latest value from Rep1);Rep2 updated in the background |
4. API Interface Design
Well-Architected: REL 3 · REL 4
The API is versioned (/v1/) so the contract can change without breaking callers (REL03-BP03), and every write accepts an idempotency key so retries are safe (REL04-BP04). PUT is already idempotent, but a retried conditional put would otherwise fail with 412 even though its first attempt succeeded.
1. Put Key-Value Record (PUT /v1/kv/{key})
Writes or updates a record. Supports idempotency tokens, conditional preconditions, and tunable write quorums.
httpPUT /v1/kv/user_profile_10928 HTTP/1.1 Host: kv.production.aws.internal Content-Type: application/json X-Consistency-Level: QUORUM X-Idempotency-Key: a8f9c1d2-7e3b-4a55-8c01-9876543210ab { "value": { "user_id": "usr_10928", "tier": "enterprise", "email": "arch@example.com", "feature_flags": ["dark_mode", "beta_billing"] }, "ttl_seconds": 2592000, "condition_expression": "attribute_not_exists(user_id) OR version = 4" }
Response: 200 OK(or 201 Createdwhen the key did not exist)
json{ "key": "user_profile_10928", "version": 5, "hlc_timestamp": "1773648000123456789-0003", "bytes_written": 184 }
2. Get Key-Value Record (GET /v1/kv/{key})
Retrieves the record with tunable read consistency.
httpGET /v1/kv/user_profile_10928 HTTP/1.1 Host: kv.production.aws.internal X-Consistency-Level: QUORUM
Response: 200 OK
json{ "key": "user_profile_10928", "value": { "user_id": "usr_10928", "tier": "enterprise", "email": "arch@example.com", "feature_flags": ["dark_mode", "beta_billing"] }, "version": 5, "hlc_timestamp": "1773648000123456789-0003", "created_at": 1773648000, "ttl_remaining_seconds": 2591985 }
For tables in sibling mode (Section 5.3), a conflicting key returns 200 OK with a siblings array instead of value; the client merges them and writes the result back.
3. Delete Key-Value Record (DELETE /v1/kv/{key})
Appends a tombstone record marking the key deleted.
httpDELETE /v1/kv/user_profile_10928 HTTP/1.1 Host: kv.production.aws.internal X-Consistency-Level: QUORUM
Response: 200 OK
json{ "key": "user_profile_10928", "deleted": true, "tombstone_version": 6, "timestamp_ns": 1773648020000000000 }
4. Batch Get Key-Value Records (POST /v1/kv:batch-get)
Retrieves up to 100 records in a single parallel scatter-gather call across different storage partitions. Returns partial results alongside unprocessed_keys if specific partitions throttle or replicas time out.
httpPOST /v1/kv:batch-get HTTP/1.1 Host: kv.production.aws.internal Content-Type: application/json X-Consistency-Level: QUORUM { "keys": [ "user_profile_10928", "user_profile_10929", "user_profile_10930" ], "projection_attributes": ["user_id", "tier", "email"] }
Response: 200 OK(Partial Success / Scatter-Gather)
json{ "records": [ { "key": "user_profile_10928", "value": { "user_id": "usr_10928", "tier": "enterprise", "email": "arch@example.com" }, "version": 5 }, { "key": "user_profile_10929", "value": { "user_id": "usr_10929", "tier": "free", "email": "free@example.com" }, "version": 2 } ], "unprocessed_keys": [ "user_profile_10930" ], "consumed_capacity_units": 4.5 }
5. Batch Write Key-Value Records (POST /v1/kv:batch-write)
Commits up to 25 items (or 16 MB) of mixed put_item and delete_item actions across multiple partition keys. The batch is not atomic: each item is committed independently on its own partition (exactly like DynamoDB BatchWriteItem), so some items can succeed while others are rejected. All-or-nothing multi-key semantics would need a separate transactional API (TransactWriteItems-style two-phase commit across partitions), which is out of scope. Items that a partition rejects because its token bucket is empty come back in unprocessed_items for the client to retry with exponential backoff, without failing the whole batch.
httpPOST /v1/kv:batch-write HTTP/1.1 Host: kv.production.aws.internal Content-Type: application/json X-Consistency-Level: QUORUM { "request_items": [ { "put_item": { "key": "user_profile_10931", "value": { "user_id": "usr_10931", "tier": "pro", "email": "pro@example.com" }, "ttl_seconds": 2592000 } }, { "delete_item": { "key": "user_profile_10932" } } ] }
Response: 207 Multi-Status/ 200 OK
json{ "committed_items": [ { "delete_item": { "key": "user_profile_10932", "deleted": true, "tombstone_version": 6, "timestamp_ns": 1773648020000000000 } } ], "unprocessed_items": [ { "put_item": { "key": "user_profile_10931", "reason": "PARTITION_THROTTLED", "retryable": true, "retry_after_ms": 25 } } ], "consumed_capacity_units": 3.0 }
Batch response payload trade-off: rich 207 Multi-Status vs sparse 200 OK
- Rich
207 Multi-Status(RFC 4918 style): returnscommitted_items(with tombstone versions and timestamps) alongsideunprocessed_items. Clients can update local caches and keep read-your-writes ordering without extra reads. - Sparse
200 OK(DynamoDBBatchWriteItemstyle): omitscommitted_itemsto cut serialization cost at very high write volume (e.g., 500k ops/sec). The client infers success by omission fromunprocessed_items. - Offer the sparse form by default and the rich form behind a request flag.
6. Architectural Deep Dive: Why PUTvs POSTin Key-Value Stores?
Choosing between HTTP PUT and POST for a key-value store follows from three distributed-systems principles:
| Dimension | PUT /v1/kv/{key} (Single Record) | POST /v1/kv (Anti-Pattern for KV) | POST /v1/kv:batch-* (Batch Actions) |
|---|---|---|---|
| Resource Identifier | Client specifies unique URI /v1/kv/{key} | Server allocates ID /v1/kv/{generated_id} | Action endpoint over many keys |
| Idempotency Guarantee | Idempotent () | Not idempotent (duplicate inserts) | Idempotent per item, not atomic (partial commit) |
| Network Retry Safety | Safe to retry on timeouts/503s | Retries create duplicate entities | Safe to retry; retry only unprocessed_items to avoid wasted capacity |
| Semantics | Upsert (create or full overwrite) | Subordinate resource creation | RPC-style scatter-gather action |
| Industry Precedent | DynamoDB PutItem, Redis SET, S3 PutObject | REST collection inserts (e.g. POST /users) | DynamoDB BatchGetItem / BatchWriteItem |
Key Technical Takeaways:
- Why not
POST(key, value)for single records?- In standard REST (RFC 9110 §9.3.3),
POSTmeans the server decides the new resource's URI. In a key-value store the client always owns and supplies the key (it is the partition key), so the resource is addressed directly at/v1/kv/{key}. PUTis idempotent. On unreliable networks a client often times out after the server already committed the write. BecausePUTis idempotent, SDKs can retryPUT /v1/kv/{key}with backoff without creating duplicates or corrupting state.- "Create-only" is a precondition, not a verb:
condition_expression: "attribute_not_exists(key)"(DynamoDB style) orIf-None-Match: *(S3 / HTTP ETag style). If the key exists, the store returns412 Precondition Failed.
- In standard REST (RFC 9110 §9.3.3),
- Why
POSTis used forbatch_getandbatch_write- URL length limits: 100 keys do not fit in
GET /v1/kv?keys=...under the 2-8 KB URL limits of reverse proxies, CDNs (CloudFront) and load balancers. - Bodies on
GET: RFC 9110 gives aGETbody no defined meaning, and many proxies and firewalls strip or reject it. - Partial failure semantics: a batch touches many partitions. If 80 keys succeed and 20 are throttled,
POSTwith200 OK+unprocessed_keys(or207 Multi-Status) models the non-atomic partial commit correctly: the verb says "an action on many resources", not "this one resource was replaced".
- URL length limits: 100 keys do not fit in
7. Status Codes & Error Contracts
| HTTP Status | Reason Code | Error Contract Payload | Client Behavior |
|---|---|---|---|
200 OK | SUCCESS | Record, or partial batch with unprocessed_* | Retry only unprocessed items |
201 Created | WRITE_COMMITTED | New key committed to replicas | None |
207 Multi-Status | PARTIAL_BATCH_COMPLETED | {"committed_items": [...], "unprocessed_items": [...]} | Back off and retry only unprocessed items |
400 Bad Request | INVALID_PAYLOAD | {"error": "Payload exceeds 1MB limit"} | Fix the request; do not retry |
404 Not Found | KEY_NOT_FOUND | {"error": "Key does not exist or tombstoned"} | Treat as absent |
409 Conflict | FENCING_TOKEN_STALE | {"error": "Request ring_epoch older than node epoch"} | Router refreshes its ring and retries internally; only token-aware SDKs (Section 7.2) see this |
412 Precondition Failed | CONDITION_CHECK_FAILED | {"error": "Condition expression evaluated to false"} | Re-read the record, then decide |
429 Too Many Requests | THROTTLED | {"error": "Partition QPS limit exceeded", "retry_after": 2} | Back off with full jitter |
503 Service Unavailable | QUORUM_NOT_REACHED / OVERLOADED | {"error": "Only 1 of 3 replicas acknowledged write"} | Back off with full jitter (honor Retry-After), then retry |
5. Data Models & Storage Architecture
1. LSM-Tree Physical Node Layout
Each storage node runs an embedded Log-Structured Merge-tree (LSM-tree) engine, the engine inside Cassandra, RocksDB and ScyllaDB, built for high-throughput concurrent writes. A write is a sequential append to the WAL plus an insert into the in-memory MemTable, so writes never do random disk I/O. Full MemTables become immutable, sorted SSTable files, and background leveled compaction merges them so each key lives in at most one file per level (below L0). A read checks the MemTable, then each level's Bloom filter, and reads one 4 KB block only from a file that may hold the key.
Each node's data directory holds three things: the commit log, the in-memory MemTable and the SSTable levels.
Path under /data/kvstore/ | What it holds |
|---|---|
commitlog/wal-0000000104.log | Append-only commit log on EBS gp3 |
memtable/ SkipList (active RAM) | 64 MB concurrent skip-list in memory |
sstables/Level-0/0001-Data.db | ZSTD-compressed data blocks |
sstables/Level-0/0001-Index.db | Sparse index (1 entry per 4 KB data block) |
sstables/Level-0/0001-Filter.db | Murmur3 Bloom filter bitset |
sstables/Level-0/0001-Summary.db | Memory-mapped index summary |
sstables/Level-1/ … Level-N/ | Leveled compaction SSTables (size ×10 per level) |
2. SSTable File Format Specification
An SSTable file is a sequence of sections, written once from the first data block to the footer:
| Section (in file order) | Size | Contents |
|---|---|---|
| Data Block 0 | 4 KB | ZSTD-compressed [KeyLen | Key | ValLen | Val | TS | Ver] |
| Data Block 1 | 4 KB | ZSTD-compressed [KeyLen | Key | ValLen | Val | TS | Ver] |
| … | … | … |
| Bloom Filter Block | Bitset array: m = 10 bits/key, k = 7 hash functions | |
| Index Block | Sparse index mapping: Key[0] → Offset 0, Key[128] → 4096 | |
| Meta Index Block | Pointers to the Bloom filter block, the stats block and the compression settings | |
| Footer | Fixed 48 B | Magic number (0x53535442), index offset, meta offset |
A reader starts at the footer, which has a fixed size, so it can always be found at the end of the file. The footer gives the offsets of the index block and the meta index block. The meta index block points to the Bloom filter, which says whether the key may be in this file at all. Only if it may be, the sparse index gives the offset of the one 4 KB data block that would hold the key, and the reader decompresses that single block.
Because SSTables never change after they are written, they are safe to cache, to share between readers without locks, and to back up incrementally (Section 10.3). Compaction amplification math and size-tiered vs leveled compaction are in Write-Ahead Log & LSM-Trees.
3. Versioning & Vector Clocks (Conflict Resolution)
Replicas can accept writes to the same key while partitioned, so every record carries a version, and each table picks one of two conflict modes:
| Mode | How a conflict is resolved | Use it for |
|---|---|---|
| Last-Write-Wins on a Hybrid Logical Clock (default) | Each write carries an HLC timestamp (physical time plus a logical counter that never goes backwards). The higher timestamp wins; the loser is discarded. | Values that are overwritten whole: profiles, sessions, flags. Simple for clients. |
| Siblings with vector clocks (opt-in) | Each write carries a vector clock {replica node: counter}, and the client sends back the clock it read (its context) on its next put. Concurrent versions are kept as siblings and returned to the client to merge. | Values where losing a concurrent update is unacceptable: shopping carts, sets, counters. |
A record in sibling mode carries its vector clock alongside the HLC timestamp:
json{ "key": "order_84920", "value": "...", "vector_clock": { "node-A1": 4, "node-B1": 2 }, "hlc_timestamp": "1773648000000000000-0000" }
- Dominance Rule: Vector clock dominates if for every node , , and for at least one node the inequality is strict.
- Concurrent Conflict: If neither dominates nor dominates , both versions are kept as siblings:
json{ "key": "cart_84920", "siblings": [ { "value": ["book"], "vector_clock": { "node-A1": 4, "node-B1": 2 } }, { "value": ["book", "pen"], "vector_clock": { "node-A1": 3, "node-B1": 3 } } ] }
Here node-A1 is higher in the first sibling and node-B1 in the second, so neither dominates; the client merges them (for a cart, the union ["book", "pen"]) and writes back with both clocks as its context.
4. Conditional Writes Need Consensus
A condition like version = 4 is a compare-and-set. Quorums alone cannot make it safe: two routers can both read version 4 from overlapping quorums and both write version 5. Conditional writes therefore run a short Paxos round among the key's 3 replicas (the same approach as Cassandra's lightweight transactions): prepare, read, propose, commit. It costs about 4 round trips instead of 1, so it is used only when the caller sends a condition. On a key that uses conditions, every write must take this path, because a plain put could overwrite a compare-and-set result by timestamp. See Distributed Consensus.
Unlock Complete Architecture & Production Runbooks
You have explored the free architectural preview (~43%). Spend 1 Coin to unlock the remaining 6 production deep-dive sections for a full 24 hours.