Hiring BarSupport

Design with Cassandra — Staff-Level Technology Guide

Technology guide43 min read5 diagrams

Why This Matters#

Cassandra is not a database you query. It is a database you pre-answer. You decide, at schema time, every question the system will ever be asked, and you lay the data on disk in exactly the order those answers will be read. Ask it a question you didn't plan for and it either refuses (ALLOW FILTERING required) or scans the cluster to find out. That's not a limitation to apologize for — it's the bargain that buys you linear write scaling, no single primary, and multi-datacenter writes that survive a region loss.

That bargain is why Cassandra keeps appearing in Staff loops. Messaging history, activity feeds, notification inboxes, IoT telemetry, fraud signals, user event timelines — any problem with high write volume, simple access patterns, and an availability requirement stronger than its consistency requirement has a Cassandra answer. The L5 candidate says "Cassandra, because it scales writes." The L6 candidate says "table notifications_by_user, partition key (user_id, month), clustering created_at DESC, RF=3 per DC, LOCAL_QUORUM both ways, TWCS with 30-day windows, TTL 90 days — and I've bounded the partition to ~2MB for a heavy user." The L7 candidate asks who will run repair for the next five years, because the most common way Cassandra deployments fail isn't load — it's operational neglect.

The L5 → L6 gap is not knowing what a memtable is. It is knowing that the partition key is the query, the partition size is the risk, and deletes are writes that come back to haunt you.

The L5 → L6 → L7 Contrast#

BehaviorSenior (L5)Staff (L6)Principal (L7)
First move"Cassandra scales writes horizontally""List the read queries first. Each one is a table; its equality predicates are the partition key.""Is the access pattern stable for three years? If product will keep inventing queries, this is the wrong engine."
Data modelNormalized tables plus a secondary indexDenormalized table-per-query, bucketed partitions sized < 10MBOwns the cost of denormalization: write amplification, multi-table consistency, and who reconciles drift
Consistency"Cassandra is eventually consistent""LOCAL_QUORUM read + write with RF=3 gives read-your-writes in-DC; LWT only for the uniqueness check"Sets which domains may use Cassandra at all: no balances, no inventory decrement, no cross-entity invariants
Deletes"Just delete old rows""Deletes are tombstones; I'll use TTL + TWCS so whole SSTables expire instead of tombstones being scanned"Bans queue-shaped workloads on shared clusters by standard, not by code review
Operations"It's self-healing""Repair inside gc_grace_seconds or deleted data resurrects; Reaper on a schedule; alert on repair age"Funds a team for it — or buys managed — because unrepaired clusters are a 2-year time bomb
Scale"Add nodes""Add nodes before 50% disk; streaming a 2TB node takes hours; plan density at 1–2TB/node"Plans cluster topology across the org: shared vs per-domain clusters, DC placement, and the exit path
Why "First move" separates levels

"Cassandra scales writes" is true and useless — so does nearly every modern distributed store. What makes Cassandra fit is that every read is a single-partition lookup or slice. If the candidate can't list the reads, they can't pick the partition key, and if they can't pick the partition key they've chosen Cassandra on reputation. The Staff move is to write the queries on the board before drawing any tables: "Q1: latest 50 notifications for a user. Q2: unread count for a user. Q3: mark notification read by ID. That's two tables and one counter — nothing scans."

Why "Deletes" separates levels

In a log-structured store, a delete is a new write — a tombstone — that must be retained long enough to reach every replica (gc_grace_seconds, 10 days by default). A read over a range that contains 50,000 tombstones has to step over each one, and at 100,000 Cassandra aborts the query. The Senior answer treats delete as free; the Staff answer designs the data so that it expires in bulk (TTL + time-windowed compaction) instead of being deleted row by row; the Principal answer keeps delete-heavy, queue-shaped workloads off Cassandra entirely as a policy.

The 60-Second Pitch#

"I'd store the notification inbox in Cassandra. It's a write-heavy, append-mostly workload — 40K inserts per second at peak, reads are always 'latest N for this user' — which maps to one partition per user per month, sorted newest first, so every read is a single-partition slice in a few milliseconds. RF=3 in each of two regions with LOCAL_QUORUM gives us read-your-writes within a region and keeps writing if a whole region goes down. Data expires by TTL at 90 days with time-window compaction, so we never pay for tombstone scans. I would not use it for anything needing cross-row transactions or ad-hoc queries — the unread-count badge is fine as a counter, but billing stays in Postgres."

The Three Intents#

IntentConstraintStrategyFailure ModeCorrectness Bar
High-volume time-ordered data (events, telemetry, messages, activity)Sustained writes, recent-first reads, retention windowPartition by entity + time bucket, cluster by time DESC, TTL + TWCSUnbounded partitions; tombstone scans if deletes replace TTLsEventually consistent; durable after LOCAL_QUORUM ack
Always-on key-value / profile store across regionsWrites accepted in every region, survive region lossNetworkTopologyStrategy, RF=3 per DC, LOCAL_QUORUM, last-write-wins per cellConcurrent cross-region updates silently resolved by timestampPer-cell LWW; must tolerate lost concurrent updates
Uniqueness or compare-and-set (usernames, idempotency keys, leases)Exactly one winnerLightweight transactions (IF NOT EXISTS) on a narrow tableLWT latency 4–10× a plain write; contention collapses throughputLinearizable per partition (Paxos)

🎯 Staff Move: "I'm choosing Cassandra for the time-ordered intent: append-heavy, read the recent slice, expire by age. The one place I need a uniqueness guarantee I'll use an LWT on a tiny dedicated table — I won't let Paxos leak into the hot write path."

The Staff Positions#

PositionRationale
Queries first, tables secondOne table per access pattern. If you can't name the query, you can't name the key.
Every partition has a size bound you can stateTarget < 10MB and < 100K rows; buckets by time or hash make "bounded" a design property, not a hope.
LOCAL_QUORUM read and write, RF=3 per DCRead-your-writes in-region, tolerates one replica loss per DC, no cross-region latency on the hot path.
Expire, don't deleteTTL + TWCS drops whole SSTables; per-row deletes create tombstones that slow reads for 10 days.
LWT is a scalpelUse it for the narrow uniqueness check, never as the default write mode.
Repair is not optionalFull or incremental repair of every token range within gc_grace_seconds, or deletes resurrect.
No read-before-write in the hot pathIt doubles latency, races anyway, and defeats the write-optimized design.

Architecture & Internals#

Only the internals that change design decisions.

The Ring, Tokens, and Replica Placement#

Every row's partition key is hashed (Murmur3) to a 64-bit token. The token space is a ring divided among nodes; each node owns several token ranges via vnodes (num_tokens, default 16 since 4.0 — down from 256, which made repair and streaming painfully slow). A partition lives on the node owning its token plus the next RF−1 nodes clockwise that sit in different racks (with NetworkTopologyStrategy), so losing one rack or AZ never removes more than one replica.

Diagram: The Ring, Tokens, and Replica Placement

There is no primary. Any node can coordinate any request. The client driver should be token-aware (send the request straight to a replica, saving a hop) and DC-aware (never cross regions for LOCAL_* levels).

The Write Path#

Diagram: The Write Path

Writes never read and never update in place: append to the commit log, insert in memory, ack. That's why a single node can absorb tens of thousands of writes per second at sub-millisecond server time — and why every update and delete is just another write that compaction must later reconcile.

Design consequences:

  • Inserts and updates are the same operation (upsert). There is no "insert fails if exists" without an LWT.
  • Last write wins per cell, by client or coordinator timestamp. Clock skew between app servers decides conflicts — use server-side timestamps or a single writer per entity when order matters.
  • Commit log sync mode (periodic every 10s by default) means a replica can lose up to ~10s of acked writes if that node crashes — which is why durability comes from RF=3 + quorum, not from fsync.

The Read Path#

A read must merge the memtable and every SSTable that might hold the partition. Bloom filters (per SSTable) skip files that definitely don't contain it; the partition index finds the offset in those that might.

Read Cost DriverWhyDesign Lever
Number of SSTables touchedEach is a potential disk seekCompaction strategy (LCS bounds it to ~1 per level)
Partition widthWide partitions mean large index + more data scannedBucket partitions; read slices with LIMIT
Tombstones in the sliceEach must be read and carried to the coordinatorTTL + TWCS; avoid range deletes in hot partitions
Consistency levelQUORUM waits for 2 of 3 replicas; digest compareLOCAL_QUORUM not QUORUM across DCs
Speculative retryCoordinator sends a redundant read when a replica is slow (default ~p99)Keeps tail latency down when one replica GCs

Compaction — the Strategy Decides Your Disk and Your Reads#

StrategyHow It WorksBest ForWatch Out For
STCS (Size-Tiered)Merge SSTables of similar sizeWrite-heavy, rarely-read, or mixed defaultNeeds up to ~50% free disk for large compactions; reads touch many SSTables
LCS (Leveled)Fixed-size SSTables in levels, ~10× per levelRead-heavy, update-heavy data~2× write amplification of STCS; struggles above ~5–10K writes/s per node
TWCS (Time-Window)Group SSTables by write-time window; never merge across windowsTTL'd time-series — whole windows drop when expiredOut-of-order writes or mixed TTLs pin old windows forever
UCS (Unified, 5.0+)Tunable to behave like STCS or LCS per level, with sharded outputNew clusters on 5.x; one strategy for most tablesNewer — validate tuning on your workload before standardizing

🎯 Staff Insight: Picking TWCS for TTL'd time-series is one of the highest-leverage single decisions in a Cassandra design. It turns "delete 2 billion expired rows" into "unlink 30 files."

Tombstones, gc_grace_seconds, and Resurrection#

A delete writes a tombstone. Tombstones must live at least gc_grace_seconds (default 864000 = 10 days) so that a replica that was down during the delete learns about it via repair before the tombstone is purged. If a replica misses the delete and repair doesn't run within that window, compaction drops the tombstone on the other replicas — and the down replica's old value comes back to life on the next read repair.

ThresholdDefaultEffect
tombstone_warn_threshold1,000 per readWarning in logs
tombstone_failure_threshold100,000 per readQuery aborted
gc_grace_seconds10 daysMinimum tombstone lifetime; your repair deadline

The Three Anti-Entropy Mechanisms#

MechanismWhen It RunsWhat It FixesWhat It Doesn't
Hinted handoffCoordinator stores writes for a down replica, replays on returnShort outages (default window max_hint_window = 3 hours)Outages longer than the window; coordinator loss
Read repairOn reads at CL > ONE where replicas disagreeHot data, opportunisticallyCold data that's never read
Anti-entropy repair (nodetool repair, Reaper)Scheduled, Merkle-tree comparison per rangeEverything, eventuallyNothing — but it's expensive and must be scheduled
Diagram: The Three Anti-Entropy Mechanisms

Data Modeling / Core Usage — "The Entire Game"#

In a relational database you model the entities and let the planner answer queries. In Cassandra you model the queries and accept that the same entity is written to several tables. Get the model right and the cluster is boring for years. Get it wrong and no amount of hardware fixes it — the only fix is a new table and a backfill.

Step 1: Write Down Every Read Query#

Worked example: a notification inbox for a 100M-user app.

#QueryFrequencyLatency Target
Q1Latest 50 notifications for user X60K/s peakp99 < 20ms
Q2Unread count for user X (badge)120K/s peakp99 < 10ms
Q3Mark notification N for user X read8K/sp99 < 20ms
Q4All notifications of campaign C (for recall)Rare, ops-onlyMinutes OK

Q4 is the trap. It has a different partition key (campaign) and is rare. Options: a second table written on every insert (write amplification for a rare query), or an offline job over a Spark/export snapshot. Rare, latency-tolerant queries belong outside the hot cluster.

Step 2: Partition Key = The Query's Equality Predicate (+ a Bound)#

Q1 is "by user, newest first." The naive key user_id puts a user's entire history in one partition — unbounded for heavy users and bots. Add a time bucket:

CREATE TABLE inbox.notifications_by_user (
    user_id     uuid,
    bucket      text,          -- 'YYYY-MM'
    created_at  timeuuid,
    notif_id    uuid,
    kind        text,
    actor_id    uuid,
    payload     text,          -- small JSON, < 2KB
    read        boolean,
    PRIMARY KEY ((user_id, bucket), created_at)
) WITH CLUSTERING ORDER BY (created_at DESC)
  AND default_time_to_live = 7776000          -- 90 days
  AND compaction = {'class': 'TimeWindowCompactionStrategy',
                    'compaction_window_unit': 'DAYS',
                    'compaction_window_size': 7}
  AND gc_grace_seconds = 172800;              -- 2 days: TTL'd, rarely deleted

Size the partition out loud:

heavy user:   200 notifications/day × 31 days = 6,200 rows/month
row size:     ~400 bytes on disk
partition:    6,200 × 400B ≈ 2.5MB   -> well under the 10MB target
pathological: a bot account at 10K/day = 310K rows, ~124MB  -> rate-limit at ingest

Q1 reads the current month's partition with LIMIT 50; if it returns fewer than 50 rows, the client reads the previous month. That's at most two single-partition reads.

Step 3: Clustering Columns = The Sort Order You Read In#

Clustering columns sort rows within a partition on disk. created_at DESC means "latest 50" is a sequential read from the start of the partition — no sort, no scan. You get one sort order per table; a second order means a second table.

Step 4: Denormalize, Then Own the Consistency Between Tables#

Q2 (unread count) can't be computed by scanning Q1's table at 120K/s. It gets its own table:

CREATE TABLE inbox.unread_counts (
    user_id uuid PRIMARY KEY,
    unread  counter
);

Now one logical event — "new notification" — writes two tables. They can diverge (a counter increment times out and is retried = double count; counters aren't idempotent). The Staff answer accepts that the badge is approximate and reconciles: when the user opens the inbox, the client recomputes unread from Q1's slice and the service resets the count. Name the owner of that drift.

Multi-Table Write OptionAtomic?CostUse When
Two independent async writesNoCheapestDerived data with a reconciliation path
Logged batch across tables (same partition key ideally)Eventually all-or-nothing (batchlog)~30–50% extra latency; coordinator loadTables must not drift, small batches (< 5KB warn threshold)
CDC/stream → derived tablesEventuallyPipeline to operateMany derived tables or other systems (search, analytics)

Step 5: Choose the Alternate-Access Strategy Deliberately#

NeedOptionVerdict
Look up by a different key, high volumeSecond table keyed by it (manual index)Default. You control size and consistency.
Filter within a partition you already knowStorage-Attached Index (SAI, 5.0+)Good — SAI on a partition-restricted query is cheap
Global lookup on a low-volume, low-cardinality columnSAI without partition restrictionOK for ops/admin queries; it fans out to every node
Legacy secondary index (2i)Avoid on high-cardinality columnsScatter-gather; latency grows with cluster size
Materialized viewsFlagged experimental; known consistency issuesAvoid in production. Maintain the second table yourself.
Arbitrary filters / full-textShip to Elasticsearch via CDCDifferent engine for a different question

Primary-Key Pattern Cheat Sheet#

WorkloadPrimary KeyWhy
User timeline / inbox((user_id, month), ts DESC)Bounded per user-month; newest first
Device telemetry((device_id, day), ts)One device-day ≈ 86,400 rows at 1Hz — fine
Chat channel history((channel_id, bucket), msg_id DESC)Bucket sized by channel activity
Global event firehose((event_type, hour, shard), ts) with shard 0–63Hash-shard to spread a hot time bucket across nodes
Profile / KV(user_id)One row, single partition
Uniqueness claim(username) + INSERT ... IF NOT EXISTSLWT on the narrowest possible table

🎯 Staff Move: "Before I write any CREATE TABLE, here are the four reads and their rates. Three of them hit Cassandra as single-partition slices. The fourth is an ops query — I'm sending it to the analytics copy so it can't hurt the inbox."

The Tunable Tradeoff — Consistency Levels#

Cassandra lets each request choose how many replicas must respond. With replication factor RF, a read at level R and a write at level W are guaranteed to overlap on at least one replica when:

R + W > RF
RF=3: QUORUM (2) + QUORUM (2) = 4 > 3   -> overlapping, read sees the latest acked write
RF=3: ONE (1) + ONE (1)       = 2 ≤ 3   -> may read a replica that hasn't seen it
RF=3: ONE write + ALL read     = 4 > 3   -> consistent, but any one node down = reads fail
LevelReplicas Waited On (RF=3/DC, 2 DCs)Typical LatencySurvivesUse For
ANY (write)Any node, even a hintLowestAlmost anythingAlmost nothing — a hint is not a replica
ONE / LOCAL_ONE1~1–3ms2 of 3 local replicas downMetrics, logs, cache-like data
LOCAL_QUORUM2 in the local DC~2–5ms1 local replica down, whole remote DC downDefault for most tables
QUORUM4 of 6 across DCs+ cross-region RTT (60–150ms)2 replicas anywhereRarely — usually a mistake in multi-DC
EACH_QUORUM (write)2 in every DC+ cross-region RTTFails if any DC is unreachableData that must be durable in every region before ack
ALL6Slowest; any node down failsNothingAlmost never
SERIAL / LOCAL_SERIALPaxos round~4 round trips (v1), ~2 (Paxos v2, 4.1+)Quorum availableLWT conditions only

What Quorum Does Not Give You#

  • Not linearizability. Two concurrent LOCAL_QUORUM writes to the same cell are resolved by timestamp, not order of arrival. If app servers' clocks differ by 50ms, the "later" write can lose.
  • Not cross-partition atomicity. A logged batch eventually applies all mutations, but readers can see a partial batch in between.
  • Not cross-DC read-your-writes. A LOCAL_QUORUM write in us-east is not guaranteed visible to a LOCAL_QUORUM read in eu-west until async replication lands (typically 100ms–1s; minutes during a partition).
  • Not "failed means not written." A write that times out at the coordinator may still have been applied on some replicas and will propagate. Retries must be idempotent.

Lightweight Transactions (LWT)#

INSERT INTO accounts.usernames (username, user_id, created_at)
VALUES ('ada', 9f2c..., toTimestamp(now()))
IF NOT EXISTS;
-- [applied] = true  -> you own it
-- [applied] = false -> returns the current owner

LWT runs Paxos on the partition: linearizable for that partition, at ~4–10× the latency of a plain write and with sharply worse throughput under contention (competing proposers retry). Mixing LWT and non-LWT writes on the same cells breaks the guarantee.

Who Pays for Each Choice#

ChoiceWho BenefitsWho Pays
LOCAL_ONE writesLatency-sensitive writersReaders who see stale/missing data; on-call explaining "lost" writes after a node crash
LOCAL_QUORUM both waysCorrectness within a region~1–2ms latency; one extra replica must be healthy
EACH_QUORUM writesCompliance/durability teamsEvery user during any cross-region blip — writes fail
LWT for every writeDevelopers who wanted a "real" databaseThroughput, p99, and the on-call during contention spikes
LWW with client timestampsSimplicityUsers whose concurrent edit silently vanishes

🎯 Staff Move: "LOCAL_QUORUM reads and writes, RF=3 per region. That gives read-your-writes within a region and keeps both regions writable if they partition. I'm accepting that a user who writes in one region and immediately reads from the other might not see it — sticky regional routing makes that rare."

Anti-Patterns — What Kills Cassandra Deployments#

1. Unbounded Partitions#

PRIMARY KEY (user_id, ts) with no bucket. Fine for two years, then a handful of power users and bots reach 1GB partitions. Reads time out, compaction of that partition stalls, repair streams the whole thing on every mismatch, and heap pressure causes GC pauses that hurt other tenants on the node. Fix requires a new table and a backfill. Prevention: a size bound in the design doc and an alert on nodetool tablehistograms max partition size > 100MB.

2. Queue-Shaped Workloads#

Insert job, worker reads oldest, deletes it. Every consumed job leaves a tombstone at the front of the partition — exactly where the next read starts. After a busy day, "read oldest" steps over 100K tombstones and fails. Cassandra is a terrible queue. Use Kafka, SQS, or a database with in-place deletes.

3. ALLOW FILTERING and Scatter-Gather Reads#

Any query without the full partition key must contact every token range. Latency grows with cluster size and one slow node sets the p99 for every such query. If it's in a request path, it's a missing table.

4. Read-Before-Write#

"Read the row, check a field, write the new value." It doubles latency, and two concurrent requests both read the old value — so it doesn't even work. Either make the write unconditional (upsert/LWW is fine), use an LWT for true conditions, or move the logic to a store with transactions.

5. Inserting Nulls#

Writing null to a column writes a cell tombstone. An ORM that sends every column on every insert can generate millions of tombstones a day without a single DELETE. Use unset values in prepared statements.

6. Large Multi-Partition Batches as a Bulk-Load Tool#

Logged batches spanning many partitions make one coordinator do all the work and write a batchlog. Batches are for atomicity across a few related writes, not throughput. For bulk loading, send many async single-partition writes or use SSTable loading.

7. Skipping Repair#

Everything works for months. Then a node that was down for 12 hours (hints expired after 3) returns, repair never ran within 10 days, tombstones are purged elsewhere — and deleted user data reappears. In a privacy context, that's a regulatory incident. Schedule repair with Reaper; alert on "time since last successful repair per table" > 7 days.

8. Oversized Nodes#

Packing 8–10TB onto a node to save money. Replacing it streams for a day or more, during which the cluster runs at reduced redundancy, and compaction can't keep up. Classic planning density is 1–2TB per node; 4.0+ zero-copy streaming and 5.0's UCS stretch that, but measure your replace time before going past ~4TB.

The Technology Landscape — Head-to-Head Comparison#

DimensionCassandraScyllaDBDynamoDBHBase / BigtablePostgres (+ Citus)
ModelWide-column, query-firstSame CQL modelKV + sort key, query-firstWide-column, sorted by row keyRelational
Write pathLSM, leaderlessLSM, leaderless, shard-per-core C++Managed, partitioned, leader per partitionLSM, single region server per key rangeB-tree + WAL, single primary per shard
ConsistencyTunable per request; LWT via PaxosTunable; LWTEventual or strong reads per request; transactions up to 100 itemsStrong per rowStrong; full ACID
Multi-region writesNative, active-active, LWWNativeGlobal Tables, LWWReplication, typically single-writerPrimary per region; hard
Ops burdenHigh: repair, compaction, JVM tuningMedium–high; fewer nodes, less GC~None; capacity/cost tuningHigh (HBase); none (Bigtable)Medium
Throughput per nodeGoodOften several × Cassandra on same hardwareN/A — per partition ~1K WCU / 3K RCUGoodBounded by primary
Ad-hoc queriesNo (SAI helps within limits)LimitedNoNo (row-key scans only)Yes
Cost shapeHardware + peopleHardware + license/peoplePay per request/storage; expensive at sustained high write ratesHardware or managedHardware
Lock-inOpen source (Apache)Licensing changed in late 2024 — check termsAWS-only APIBigtable: GCP-only; HBase: openOpen source

How to choose in one breath: sustained high write volume, multi-region active-active, known queries, and a team willing to operate it → Cassandra (or ScyllaDB for density). Same workload but no appetite to run a cluster and a single cloud → DynamoDB. Need joins, transactions, or ad-hoc queries → Postgres, sharded when you must. See DynamoDB, PostgreSQL, and the Database Selection case study.

Patterns#

Pattern 1: Time-Bucketed Event Store#

The canonical use. Partition by (entity, time_bucket), cluster by time, TTL the data, TWCS with windows ≈ 1/20th–1/30th of the TTL (e.g., 90-day TTL → 3–7 day windows; keep total windows under ~50). Pick the bucket so the busiest entity's partition stays < 10MB.

bucket = floor(ts / bucket_width)
bucket_width chosen so: max_events_per_entity_per_bucket × row_bytes < 10MB
reads: newest bucket first, walk back until LIMIT satisfied (cap the walk at N buckets)

Pattern 2: Hash-Sharded Hot Buckets#

When the partition key is shared by everything in a time window (a global feed, "all events this hour"), one partition takes the whole write load. Add a synthetic shard: ((hour, shard), ts) with shard = hash(event_id) % 64. Writes spread over 64 partitions; reads fan out 64 single-partition queries in parallel and merge. Choose the shard count from peak write rate ÷ comfortable per-partition rate (~a few thousand writes/s).

Pattern 3: Manual Index Table#

Need to find a notification by ID (Q3, "mark read")? Write notifications_by_id (notif_id → user_id, bucket, created_at) alongside the main row, then update the main row by its full key. Two writes per insert; reads stay single-partition.

Pattern 4: Multi-DC Active-Active#

Diagram: Pattern 4: Multi-DC Active-Active

Each region reads and writes locally; replication is asynchronous; conflicts resolve by last-write-wins per cell. Sticky routing by user keeps a single user's writes in one region most of the time, which makes LWW conflicts rare in practice. On region loss, the global LB shifts traffic; the surviving DC already has (almost) all data. After recovery, run repair on the returning DC before trusting it for reads at lower levels.

Pattern 5: Uniqueness via LWT Side Table#

Keep the main table LWT-free. Claim the unique value in a tiny table with IF NOT EXISTS, then write the main row unconditionally. If the main write fails, a sweeper (or a TTL on unconfirmed claims) releases orphans.

Pattern 6: CDC Out of Cassandra#

Cassandra's CDC writes commit-log segments to a CDC directory per node; consuming it means de-duplicating across RF replicas. Many teams instead emit events from the application (outbox-style) into Kafka, and feed search, analytics, and derived tables from there. Use CDC when the application can't be changed.

Scaling#

The Numbers#

MetricPlanning ValueNotes
Writes per node~10–25K replica writes/s (8–16 cores)Client writes/s × RF = replica writes/s
Reads per node~5–15K single-partition reads/sDepends heavily on SSTables-per-read and cache hits
Write latency (LOCAL_QUORUM)p50 ~1–2ms, p99 ~5–10msSame AZ set
Read latency (LOCAL_QUORUM)p50 ~2–4ms, p99 ~10–20msWide partitions and tombstones dominate p99
Data per node1–2TB classic; ~4TB with care on 4.x/5.xLimited by streaming/replace time, not by disk
Partition size< 10MB target, < 100MB hard ceilingRows < 100K per partition
Cell value< 1MB, ideally < 100KBLarge blobs → object storage + pointer
Free disk headroom30–50% (STCS needs most)Compaction needs room to write the merged SSTable
num_tokens16 (4.0+)256 vnodes made repair and streaming slow
Heap8–31GB (G1/ZGC); keep < 32GB for compressed oopsOff-heap for bloom filters, memtables optional

Sizing Example#

Ingest:      40K notifications/s peak, 400 bytes each, RF=3, 2 DCs
Writes/DC:   40K × 3 = 120K replica writes/s  -> at ~15K/node = 8 nodes for writes
Storage:     avg 15K/s (peak 40K) × 86,400 × 90 days × 400B ≈ 47TB raw per replica set
Per DC:      47TB × 3 (RF) ≈ 140TB; with 40% headroom ÷ 0.6 ≈ 233TB
Nodes/DC:    233TB ÷ 2TB/node ≈ 117 nodes  -> storage, not throughput, is the driver
Lever:       compression (LZ4 ~2–3× on JSON payloads) -> ~45–60 nodes/DC

Storage, not request rate, sizes most Cassandra clusters. The two cheapest levers are compression and a shorter TTL — the second is a product conversation.

Scaling Moves in Order#

  1. Fix the model — wide partitions and tombstones cause more "we need more nodes" than load does.
  2. Tune compaction — wrong strategy can triple read latency or disk use.
  3. Add nodes — one at a time (or with care, a rack at a time); each bootstrap streams ~1/N of the data; run nodetool cleanup on existing nodes afterward.
  4. Split the cluster — separate workloads with different shapes (time-series vs KV) onto different clusters so compaction and GC don't fight.
  5. Add a DC — for latency or region resilience; ALTER KEYSPACE RF for the new DC, then nodetool rebuild from an existing DC.

Failure Modes & Recovery#

1. Tombstone Overwhelm#

  • Symptom: a specific query starts timing out; logs show Read N live rows and M tombstone cells; eventually TombstoneOverwhelmingException.
  • Root cause: row-level deletes or null inserts concentrated where reads start (queue pattern, "delete read notifications").
  • Detection: cassandra.table.tombstone_scanned_histogram p99 > 1,000; warn-log rate.
  • Fix: short term, lower gc_grace_seconds on that table only if repair is current, then force compaction; long term, redesign to TTL + TWCS or move the workload.
  • Prevention: design review rule — no delete-heavy read-front patterns; ORM unset-values check.

2. Hot or Wide Partition#

  • Symptom: one node's CPU and latency spike while peers idle; timeouts for a subset of users.
  • Root cause: a celebrity, bot, or global key concentrating reads/writes; or an unbounded partition grown to hundreds of MB.
  • Detection: nodetool tablehistograms partition size max; per-node ClientRequest.Latency skew; driver-side per-host latency.
  • Fix: rate-limit the abusive key at ingest; add a hash shard to the key (new table + backfill).
  • Prevention: size bound in every design doc; alert on max partition > 100MB.

3. Compaction Backlog → Disk Full#

  • Symptom: PendingCompactions climbs for hours, SSTable count per table grows, read latency rises, disk usage approaches 90%.
  • Root cause: sustained write burst beyond compaction throughput, or STCS major compaction needing ~50% free space that isn't there.
  • Detection: cassandra.compaction.pending_tasks > 100 for 30 min; disk > 70%.
  • Fix: raise compaction_throughput temporarily, add nodes (streaming also needs disk — act early), drop expired data via TTL, delete snapshots (a common hidden disk eater).
  • Prevention: alert at 60–70% disk, not 90%; capacity plan for compaction headroom.

4. Data Resurrection After Missed Repair#

  • Symptom: deleted rows (account deletions, GDPR erasures) reappear weeks later.
  • Root cause: a replica missed the delete, hints expired (> 3h outage), repair didn't complete within gc_grace_seconds, tombstones were purged elsewhere.
  • Detection: repair.last_success_age per table > gc_grace_seconds × 0.7; audit jobs comparing erasure log to live data.
  • Fix: re-issue deletes from the erasure log; run full repair; never bring a node back that has been down longer than gc_grace_seconds — wipe and replace it instead.
  • Prevention: Reaper with per-table schedules; hard runbook rule on long-down nodes.

5. Coordinator Overload and GC Pauses#

  • Symptom: p99 spikes cluster-wide in waves; dropped mutations; nodes flapping DOWN/UP in gossip.
  • Root cause: large partitions or big batches materialized on heap; un-token-aware clients making every node coordinate; heap too small for the workload.
  • Detection: jvm.gc.pause > 500ms; DroppedMessage.MUTATION > 0; ThreadPool.Native-Transport-Requests.pending.
  • Fix: enable token-aware routing in drivers; cap page size (fetch_size ~1–5K rows); move to G1/ZGC and tune; find the big partition.
  • Prevention: driver configuration as a shared library default; GC pause SLO per node.

Operational Reality Matrix#

FailureDetection SignalBlast RadiusMitigationOwner
Tombstone overwhelmtombstone_scanned_histogram p99One table's queriesTTL/TWCS redesign; targeted compactionOwning service team
Hot / wide partitionPer-node latency skew, partition size maxAll tables on affected replicasIngest rate limit, re-keyService team (model) + DB platform (detection)
Compaction backlog / diskpending_tasks, disk %Node → cluster if replace neededThroughput tuning, add nodes, clear snapshotsDB platform
ResurrectionRepair age vs gc_grace_secondsCorrectness, complianceRe-delete, full repair, replace long-down nodesDB platform; privacy team audits
GC / coordinator overloadGC pause, dropped mutationsCluster-wide p99Token-aware drivers, page sizes, heap tuningDB platform + client library owners
DC lossGossip: DC nodes DOWNOne region's local quorumsShift traffic; rebuild/repair on returnSRE / DB platform

When to Use vs. Alternatives#

RequirementPickWhy
> 50K writes/s sustained, known queries, multi-region writesCassandra / ScyllaDBLeaderless, linear write scaling, active-active
Same, but no team to run it, single cloudDynamoDB (or a managed Cassandra service)No repair, compaction, or JVM to own
Joins, transactions, evolving queriesPostgresPlanner + ACID; shard later
Balances, inventory, ledgerPostgres / Spanner-class / CockroachDBCross-row invariants need real transactions
Metrics with downsampling and aggregationTime-series DBPurpose-built compression and rollups — see Time Series DBs
Full-text, faceted searchElasticsearch fed from the source of truthDifferent index structure — see Elasticsearch
Work queueKafka / SQSTombstones make Cassandra a bad queue

When NOT to Use Cassandra#

  • Your queries aren't known yet. Early-stage products change access patterns monthly; each change is a new table and a backfill.
  • You need cross-entity invariants. "Don't oversell," "balance ≥ 0," "exactly one active session" — LWT handles one partition; anything wider is application-level hope.
  • Data is small. Under ~1TB and ~10K writes/s, a well-run Postgres primary with replicas is cheaper and simpler.
  • Nobody will own operations. Repair, compaction tuning, node replacement and upgrades are real work. If no one is funded for it, buy managed or pick a different engine.
  • The workload is delete-heavy or queue-like. Tombstones will win.

Operational Concerns#

What the On-Call Actually Does#

  • Replaces failed nodes (replace_address_first_boot) and watches streaming finish; ~1–4 hours for a 1–2TB node on 10Gbps.
  • Watches repair progress in Reaper; restarts failed segments.
  • Responds to disk alerts: clears snapshots, checks pending compactions, decides whether to add capacity now.
  • Finds the wide partition or tombstone table behind a latency page, then hands the fix to the owning team.
  • Runs rolling restarts for config, JVM, and OS patches — one node at a time, waiting for UN status and hint replay.

Key Metrics & Alerts#

MetricAlert ThresholdWhy
ClientRequest.Read/Write.Latency p99 per DC> 2× baseline for 10 minUser-visible slowness
ClientRequest.Unavailables / Timeouts> 0 sustainedQuorum not reachable
compaction.pending_tasks> 100 for 30 minBacklog → read amplification → disk
Disk used %> 60% warn, > 75% pageCompaction and streaming need room
Repair age per table> 70% of gc_grace_secondsResurrection risk
Max partition size> 100MBModel problem in progress
DroppedMessage.MUTATION> 0Writes shed under load — hints/repair must catch up
JVM GC pause> 500msCoordinator stalls, gossip flaps
Hints in flightGrowing for > 1hA replica is down or slow

Upgrades and Config Changes#

Rolling, one node at a time, in one DC first; never mix major versions longer than the upgrade window; don't run repair or schema changes mid-upgrade; run upgradesstables after a major version jump. Schema changes (ALTER TABLE) propagate via gossip — make them from one place, one at a time, and wait for schema agreement; concurrent schema changes from multiple services are a classic self-inflicted outage.

Interview Application — Staff-Level Plays#

Which Case Studies Use Cassandra#

Case StudyWhere Cassandra FitsThe Staff Detail
Chat MessagingMessage history per conversation((conversation_id, bucket), message_id DESC); bucket sized by the busiest channel
News FeedPrecomputed timelines, activity storeFan-out-on-write into timeline_by_user; TTL old entries; celebrities handled on read
Notification SystemInbox and delivery logTTL + TWCS; unread counter is approximate and reconciled
Metrics & MonitoringRaw sample storage (historically)Usually lose to purpose-built TSDBs — say why
Ad Click AggregatorAggregated counts by campaign/timeWrite idempotent aggregates, not counters, so replays are safe
Web CrawlerURL frontier metadata, crawl historyPer-host partitions; avoid the queue anti-pattern for the frontier itself
Replicated Data StoreReference design for leaderless quorum replicationR + W > RF, hinted handoff, read repair, Merkle-tree repair
Database SelectionThe "write-heavy, known queries, multi-region" branchContrast with DynamoDB and Postgres on ops cost

Every System Design Question Has a Cassandra Moment#

  • Chat: "Messages go in Cassandra keyed by conversation and time bucket — append-only, read newest-first, and a bucket bounds the partition even for a 50K-member channel."
  • Feed: "Timelines are a materialized answer per user. That's exactly what Cassandra is good at — one partition read per feed load."
  • IoT / telemetry: "Device-day partitions, TTL at retention, TWCS — old data disappears as whole files, never as tombstones."
  • Anything with money: "Not here. Balances need cross-row invariants; that stays in a transactional store and Cassandra gets the event history."

What Interviewers Probe#

After You Say...They Will Ask...What They're Evaluating
"Cassandra for the messages""What's the partition key? How big does the largest partition get?"Whether you size partitions or hope
"LOCAL_QUORUM""A user writes in Europe and reads in the US. What do they see?"Understanding of async cross-DC replication
"We delete read notifications""What happens to read latency a week later?"Tombstone awareness
"We'll add a secondary index for that""How does that query execute on a 60-node cluster?"Scatter-gather awareness; table-per-query discipline
"Usernames must be unique""How, without a primary?"LWT and its cost; side-table pattern
"Cassandra is self-healing""A node was down for two weeks. Do you bring it back?"gc_grace_seconds and resurrection
"It scales linearly""What does adding a node actually cost, operationally?"Streaming, cleanup, density limits

Common Interview Mistakes#

What Candidates SayWhat Interviewers HearWhat Staff Engineers Say
"Cassandra because it's fast"Chose by reputation"Cassandra because the writes are 40K/s, append-only, multi-region, and every read is one partition"
"Partition key is user_id"Hasn't sized the worst user"(user_id, month) — the heaviest user is ~2.5MB per partition"
"It's eventually consistent, so it's fine"Hand-waving correctness"LOCAL_QUORUM both ways gives read-your-writes in-region; cross-region lag is ~1s, and sticky routing hides it"
"We'll query by status with ALLOW FILTERING"Will page someone at scale"That's a second table keyed by status and bucket — or it's an analytics query"
"Delete expired rows nightly"Will create a tombstone problem"TTL at write, TWCS windows, so expiry drops whole SSTables"
"Use counters for the click totals"Will double-count on retries"Counters aren't idempotent. I'll write per-window aggregates keyed by window and overwrite them"

L5 vs L6 vs L7 Responses#

ScenarioL5 AnswerL6 / Staff AnswerL7 / Principal Answer
"Store chat messages""Cassandra with channel_id partition key""((channel_id, bucket), msg_id DESC); bucket width chosen so the busiest channel stays < 10MB; LOCAL_QUORUM; edits are upserts, deletes are soft flags""Message history is a 10-year retention asset. I'll decide hot vs cold tiering now — recent months in Cassandra, older in object storage with an index — so node count doesn't grow forever"
"A query is slow""Add a secondary index""Which table should answer that query? Build it, dual-write, backfill, cut over""Why did a new query reach the hot cluster without a design review? Add a data-model review gate for shared clusters"
"We need strong consistency for one field""Use QUORUM everywhere""LWT on that one partition, or move that field to a transactional store""Define the org rule: which invariants are allowed on leaderless stores at all"
"Expand to a second region""Add a DC and replicate""NetworkTopologyStrategy RF=3 in new DC, nodetool rebuild, LOCAL_QUORUM, sticky user routing, LWW conflict analysis""Price cross-region replication bandwidth and decide which keyspaces actually need to be in both regions — not all do"
"The cluster is expensive""Use bigger nodes""Compression, shorter TTL, right compaction strategy, then density""Compare 3-year TCO of self-run vs managed vs ScyllaDB, including people and the migration cost"

The Staff Cassandra Checklist#

  1. "Here are the read queries and their rates — each gets a table."
  2. "Partition key is the equality predicate plus a bucket; here's the worst-case partition size."
  3. "Clustering order matches how we read — newest first, sliced with LIMIT."
  4. "RF=3 per DC, LOCAL_QUORUM reads and writes; here's what that means across regions."
  5. "Data expires by TTL with TWCS; we don't delete row-by-row."
  6. "Repair runs on a schedule inside gc_grace_seconds, and the platform team owns it."

🎯 Staff Insight: Say what you won't put in Cassandra. "Balances, inventory counts, and anything that needs a uniqueness or ordering guarantee across rows stay in a transactional store. Cassandra gets the high-volume history." The boundary is the signal.

Evaluation Rubric#

DimensionSenior (L5)Staff (L6)Principal (L7)
Data modelingCorrect partition key for the main queryTable per query, bounded partitions with size math, clustering order matches readsGoverns model changes on shared clusters; plans hot/cold tiering for multi-year retention
ConsistencyKnows QUORUM and ONER + W > RF, LOCAL_QUORUM, LWT boundaries, LWW pitfallsSets org policy on which invariants may live on leaderless stores
Deletes & lifecycle"Delete old data"TTL + TWCS, tombstone thresholds, gc_grace_secondsMakes retention a product decision with a cost attached
Operations"Add nodes"Repair, compaction strategy, node density, replace times, alertsFunds/staffs operations or buys managed; owns the exit plan
FitPicks Cassandra for scalePicks it for the write pattern and rejects it for transactional partsDecides whether the org should run it at all, and for which domains

Strong Hire Signals

SignalWhat It Sounds Like
Queries before tables"Four reads. Three are single-partition; the fourth is an analytics job."
Sizes the worst partition"Heaviest user at 200/day is 2.5MB per month bucket."
Designs for expiry"TTL with 7-day TWCS windows — expiry is file deletion."
Draws the transactional boundary"Counts and balances don't live here."
Owns the repair"Reaper weekly, alert at 70% of gc_grace, platform team on call."

Lean No-Hire Signals

SignalWhy It Misses the Bar
Normalized schema with joins "in the app"Turns every read into N partition reads
Secondary indexes for primary access pathsScatter-gather in the hot path
"Eventually consistent" with no levels or numbersCan't reason about what users see
Row-level deletes for retentionTombstone failure waiting to happen
No mention of operationsPicked a database they couldn't run

Common False Positives

  • Knowing SSTable, memtable, and bloom-filter internals ≠ modeling data correctly.
  • Saying "Netflix uses it" ≠ justifying it for this workload.
  • Choosing QUORUM everywhere ≠ understanding consistency — in multi-DC it's usually the wrong level.

The Principal Lens#

Why L7 Sees This Problem Differently#

At Staff level, Cassandra is an engine you model correctly. At Principal level it is a long-term operational commitment with a skill requirement most orgs underestimate. The design that's perfect on day one is still running in year five, with three times the data, two reorgs later, owned by a team that didn't build it. The failures that hurt at that point are not modeling mistakes — they are skipped repairs, a node density nobody re-evaluated, and a cluster shared by six teams where one team's tombstones degrade everyone. The L7 question is "Should we be in the business of operating a leaderless database — and if so, how many clusters, for whom, with what guardrails, and what's our exit?"

The Org-Level Fault Line#

One shared Cassandra platform vs. per-domain clusters vs. buying it.

OptionWhat WorksWhat BreaksWho Pays
One large shared clusterUtilization; one expert team; one upgrade cycleNoisy neighbors (one team's wide partition = everyone's GC pauses); schema changes collideEvery tenant on the bad day
Per-domain clusters from a platform (messaging, telemetry, profiles)Isolation; per-workload compaction/heap tuning; independent upgradesMore clusters to patch; idle capacityPlatform headcount (~1 FTE per 8–15 clusters with good automation)
Teams self-runAutonomyRepair silently stops; inconsistent versions; resurrection bugsPrivacy/compliance, incident responders
Managed (Keyspaces, Astra, Azure Managed Instance) or DynamoDBNo repair/compaction opsFeature gaps vs Apache Cassandra; per-request pricing at high write volume; vendor lock-inFinance; teams adapting to feature differences

The Principal default: per-domain clusters provisioned and operated by one data-platform team, with a data-model review gate for any new table on a shared cluster, automated repair as part of the platform, and a written decision for which domains should use a managed service instead.

Cost Model#

Assumptions: self-run on cloud VMs with local NVMe, ~$700–1,000/node-month for an 8–16 vCPU / 64GB / 2TB node; 2 DCs; cross-region transfer ~$0.02/GB; loaded engineer ~$250K/yr. Managed estimates assume sustained write-heavy usage. Directional only.

ScaleData (per DC, post-RF)Self-Run Infra/monthCross-Region/monthPeopleManaged Estimate/month
Startup3TB, 5K writes/s~$2.5–4K (3 smaller nodes × 2 DCs)~$0.5–1K0.5 FTE (~$10K) — the real cost~$5–15K
Growth150TB, 100K writes/s~$70–100K (~100 nodes/DC)~$10–25K2–3 FTE (~$50K)Often 1.5–3× self-run at this write rate
Enterprise2PB, 1M+ writes/s, 3 DCs~$1–1.5M (1,000+ nodes)~$100K+8–15 FTE + toolingRarely viable; negotiate or self-run

At startup scale, people dominate — half an engineer costs more than the cluster, which is the strongest argument for a managed or different engine. At growth scale, storage density and compression dominate the node count. At enterprise scale, cross-region replication and TTL policy become product-level cost decisions.

The 3-Year Evolution Path#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionReversibilityCost to Reverse
Partition key of a large tableOne-wayNew table, dual-write, full backfill, cutover — weeks to months at TB scale
Choosing a leaderless LWW store for data that needs invariantsOne-way in practiceData cleanup after anomalies, then migration to a transactional store
Retention promise to customers ("we keep 5 years of history")One-wayContract changes; storage grows forever
Which keyspaces replicate to which regionsOne-way-ish (residency)Legal review; data deletion verification
Compaction strategyTwo-wayALTER TABLE; rewrite cost on large tables
Node size / instance typeTwo-wayRolling replace
Self-run → managed Cassandra-compatibleTwo-way-ishCQL is portable; features and limits differ — test first

🧭 Principal Move: "Teams can tune compaction, TTLs and node sizes without asking anyone. Partition keys on shared clusters and anything that replicates personal data across regions go through review — those are the decisions we pay for for years."

The Standard I'd Write#

RFC-DATA-014: Using Cassandra on the Data Platform

Scope: Any table on a platform-operated Cassandra cluster.

MUST
  1. Submit the query list, partition key, and worst-case partition size
     (< 100MB, target < 10MB) for review before table creation.
  2. Use LOCAL_QUORUM (or justify otherwise) and token-aware, DC-aware drivers
     from the platform client.
  3. Express retention as TTL; tables with > 1% of writes being deletes need an
     approved exception.
  4. Never store data requiring cross-partition invariants (balances, inventory,
     unique constraints beyond a single LWT partition).
  5. Platform guarantees repair of every table within 70% of gc_grace_seconds and
     alerts on breach.
SHOULD
  6. Use TWCS (or UCS) for TTL'd time-series; LCS only for read-heavy update data.
  7. Tier data older than the hot window to object storage.
  8. Emit domain events via outbox to Kafka rather than consuming Cassandra CDC.

Exceptions: Approved by data-platform lead and the requesting team's director;
reviewed every two quarters.

Success metrics: zero resurrection incidents; zero tables with partitions > 100MB;
p99 read < 20ms per cluster; storage cost per retained TB trending down.

What I'd Tell the VP#

"Cassandra is the right engine for our message and activity history: it handles our write volume across two regions and keeps working if one goes down. The risk isn't the technology, it's that it needs steady care — maintenance jobs that, if skipped, can bring deleted customer data back. Today that work depends on two people. I'm proposing we make it a platform responsibility with automation, a review step for new tables, and a rule that financial data never goes into it. That's roughly one additional engineer and protects both our uptime and our privacy commitments. We'll also revisit, next year, whether a managed service is cheaper than running it ourselves."

Principal Interview Signals#

SignalWhat It Sounds Like
Treats operations as the long-term cost"The cluster is $3K a month; the half-engineer to run it is $10K. That's the real decision."
Gates one-way doors"Partition keys get reviewed; compaction settings don't."
Separates workloads by blast radius"Telemetry and profiles don't share a heap."
Thinks in retention economics"Every year of retention is ~40 nodes. Product should choose that knowingly."
Plans the exit"CQL compatibility gives us managed and alternative-engine options; we test that path annually."

Staff answers that L7 interviewers find insufficient:

  • "Run repair weekly" — without saying who owns it, how breaches are detected, or what happens when that person leaves.
  • "Use a separate table for that query" — correct, but silent on who reviews the next ten tables other teams add to the same cluster.
  • "Cassandra is cheaper than DynamoDB at this scale" — without the people cost in the comparison.

In the Wild#

Apple — Scale as an Operational Discipline#

Apple has spoken publicly at Cassandra conferences about running one of the largest deployments in existence — well over 100,000 nodes storing many petabytes — and Apple engineers are among the significant contributors to the Apache project. At that size, the dominant engineering problems are automation: node replacement, repair, upgrades, and fleet-wide configuration.

Staff insight: The largest users invest most in operating Cassandra, not in modeling. In an interview, mentioning who runs repair and replacement is what distinguishes "I've read about it" from "I've carried the pager."

Netflix — Multi-Region, Active-Active#

Netflix adopted Cassandra during its move to AWS for data that needed to stay writable across availability zones and regions, published early benchmarks showing over a million writes per second on a few hundred nodes, and built open-source tooling (such as Priam for backups and token management) around it. Later, Netflix described building a data-access abstraction layer in front of its key-value stores so application teams program against a stable API rather than raw CQL.

Staff insight: The abstraction layer is the Principal lesson: once dozens of teams use a store, the org wants a contract it controls between applications and the engine, so data models and engines can evolve without touching every service.

Discord — Partition Design and Its Limits#

Discord's 2017 engineering post described storing messages in Cassandra with a partition key of channel plus a fixed time bucket, so even very active channels had bounded partitions. Their 2023 follow-up described the problems that emerged at trillions of messages — hot partitions in huge servers, GC pauses, and heavy maintenance — and a migration to ScyllaDB fronted by data services that coalesce concurrent reads of the same hot data, with a substantial reduction in node count.

Staff insight: Both halves are interview gold. The first shows the bucketed-partition pattern done right; the second shows that even a correct model hits hot-partition and operational limits at extreme scale — and that request coalescing in front of the database is often the fix, not just a new engine.

Practice Drill#

Prompt: "Design storage for a connected-car platform: 2 million vehicles each send a 300-byte telemetry record every 5 seconds while driving (assume 20% are driving at peak). Fleet managers view the last 24 hours of a vehicle's trip on a map; data science wants 13 months of history; a vehicle's owner can request deletion of all their data."

Staff Answer

Peak ingest: 2M × 20% ÷ 5s = 80K writes/s, ~24MB/s raw. Reads: "last 24h for vehicle V" — a single-entity, recent-first time slice — so Cassandra fits, keyed ((vehicle_id, day), ts DESC). Partition size: worst case a vehicle driving 24h = 17,280 rows × ~300B ≈ 5MB — under the 10MB target, so a day bucket works; the map view reads today's and yesterday's partitions. RF=3 in each of two regions, LOCAL_QUORUM writes from regional ingest, LOCAL_ONE acceptable for the fleet map (a second of staleness on a trip trace is fine) — I'll still use LOCAL_QUORUM by default and relax only if p99 demands. Retention split: the hot table keeps 30 days with default_time_to_live = 30 days and TWCS 1-day windows (~30 windows). The 13-month history doesn't belong on the hot cluster: an ingest consumer writes the same stream from Kafka to Parquet in object storage, partitioned by day and vehicle hash, which data science queries with Spark/Trino — roughly 10× cheaper per TB than 3× replicated NVMe. Sizing the hot tier: average maybe 25K writes/s × 86,400 × 30 days × 300B ≈ 19TB per replica set, × 3 RF ≈ 58TB per DC, ~2–3× compression → ~25TB, ÷ 0.6 headroom ÷ 2TB/node ≈ 20–24 nodes per DC; throughput check: 80K × 3 = 240K replica writes/s ÷ ~15K/node ≈ 16 nodes — storage and throughput agree around 24. Deletion: a user-deletion request doesn't delete rows one by one across 30 day-partitions — it issues a partition-level delete per (vehicle_id, day) (30 partition tombstones, cheap) and a deletion job against the object-store history; the platform's repair SLO (complete within 70% of gc_grace_seconds, set to 3 days on this TTL'd table) is what makes the deletion stick. Owners: telematics team owns the model and retention; data platform owns cluster, repair, and the lake export; privacy team audits deletion completion.

Why this is L6:

  • Derives write rate, partition size, and node count from the prompt's numbers, and checks storage against throughput.
  • Splits hot (Cassandra, 30 days) from cold (object storage, 13 months) instead of retaining everything in the expensive tier.
  • Uses TTL + TWCS for retention and partition-level deletes for erasure — no row-level tombstone storm.
  • Connects deletion correctness to repair within gc_grace_seconds and names owners.

What L7 adds:

  • Makes 30-day hot retention a product decision with a price: each extra hot month is ~20 nodes per region.
  • Requires the export to the lake to be a platform pipeline shared with other telemetry domains, not a one-off consumer.
  • Puts vehicle-data residency into the keyspace design up front — EU vehicles' keyspace replicated only within the EU — because that is a one-way door once data has crossed regions.

Quick Reference Card#

Model:            queries first; one table per access pattern
Primary key:      ((partition cols), clustering cols) -> partition = query's equality predicate
Partition size:   target < 10MB, ceiling 100MB, < 100K rows; bucket by time or hash
Consistency:      R + W > RF;  default RF=3/DC + LOCAL_QUORUM read & write
Cross-DC:         async, ~100ms–1s; LWW per cell by timestamp
LWT:              IF NOT EXISTS / IF col = x; Paxos; ~4–10x write latency; narrow tables only
Writes:           upserts; commit log + memtable; no read-before-write
Deletes:          tombstones; live >= gc_grace_seconds (default 10 days)
Tombstone limits: warn 1,000 / fail 100,000 per read
Repair:           every range within gc_grace_seconds (Reaper); hints window 3h
Compaction:       TWCS for TTL time-series; LCS read-heavy; STCS write-heavy; UCS on 5.x
Node density:     1–2TB classic, ~4TB with care; 30–50% disk headroom
vnodes:           num_tokens = 16 (4.0+)
Indexes:          manual index tables > SAI (partition-restricted) > 2i; avoid MVs
Throughput:       ~10–25K replica writes/s per node; storage usually sizes the cluster

RED FLAGS
  - Partition key with no bound (user_id alone for an event stream)
  - Queue / delete-heavy workload
  - ALLOW FILTERING or 2i in a request path
  - QUORUM (not LOCAL_QUORUM) in multi-DC
  - Counters for anything that's retried
  - No repair schedule, or a node returning after > gc_grace_seconds
  - Balances or inventory in a LWW store
  1. Loading the index…