Cassandra & ScyllaDB Distributed Architecture Studio
Calculate partition sizes and cell counts, model tunable consistency quorum math (R + W > N), select SSTable compaction strategies, and synthesize production CQL DDL.
Partition Dimensionality & Sizing Modeling
Partitions exceeding 100MB or 100,000 cells trigger severe JVM GC pauses and CPU core imbalance.
Under the 100MB limit (14.8 MB) and under 100k cells (43,200 cells). Safe for production reads and SSTable compactions.
• Partition Key Cache: Easily cached in Linux page cache and row cache.
• Bucketing Advice: If partition exceeds 100MB, add a time bucket to the partition key (e.g.
PRIMARY KEY ((device_id, yyyymmdd), timestamp)).
Tunable Consistency & Quorum Math
Linearizable consistency requires $R + W > N$. Configure Read and Write consistency levels and test cluster node failure tolerance.
Because $R + W > N$, at least one replica node is guaranteed to participate in both the write and read operations, ensuring stale data is never returned.
SSTable Compaction Strategy Decision Guide
Selecting the wrong compaction strategy leads to 50% wasted disk headroom, high read latencies, or catastrophic write amplification.
• Workload: Write-heavy append logs with few reads.
• Mechanics: Groups SSTables of similar sizes (4 SSTables of 100MB $
ightarrow$ 1 SSTable of 400MB).
• FATAL DRAWBACK: Requires 50% free disk space during major compactions! If disk crosses 50% utilization, compactions stall and node runs out of disk.
• Workload: Read-heavy workloads (e.g. 90% reads, 10% writes, user profiles).
• Mechanics: Organizes SSTables into hierarchical levels (L0, L1, L2). Each level is 10x larger than previous.
• Advantage: Guarantees 90% of reads hit at most 1 SSTable. Requires only 10% free disk space.
• Drawback: High write amplification (rewrites data ~10 times).
• Workload: Time-series, IoT telemetry, event streams with TTL.
• Mechanics: Partitions SSTables into discrete time windows (e.g. 1 hour or 1 day). Once a time window closes, compactions cease forever.
• Advantage: When TTL expires, Cassandra purges the entire SSTable file immediately with zero disk read/write overhead!
5 Architectural Showdowns & Decision Matrices
Cassandra runs on Java HotSpot JVM; heavy garbage collection pauses (Stop-The-World) can cause nodes to drop heartbeats and trigger read timeouts. ScyllaDB rewrote Cassandra in C++ on the Seastar asynchronous reactor framework with thread-per-core architecture, delivering 5x higher throughput with consistent sub-millisecond p99 latencies.
Cassandra: Self-hosted masterless peer-to-peer ring; zero single point of failure, no vendor lock-in. DynamoDB: Fully managed AWS serverless, but subject to 1,000 WCU/RCU partition throttle limits. MongoDB: Document database with single primary write replica; master failover takes 2–10 seconds.
In multi-DC setups, QUORUM requires majority across the entire world, stalling local requests on trans-oceanic WAN roundtrips. LOCAL_QUORUM confines coordination strictly to local DC nodes, surviving inter-DC cable cuts with zero latency penalty.
Cassandra native Materialized Views execute asynchronous shadow mutations in the background; during cluster overload, view updates lag indefinitely without backpressure. Production enterprise systems implement Application Dual-Writes or Kafka consumer sinks for predictable secondary indexing.
5 Fatal Production Cassandra Pitfalls
Writing null columns or deleting records creates tombstones. When a single read query traverses >100,000 tombstones, Cassandra terminates the query with
TombstoneOverwhelmingException. Remedy: Never insert nulls, use TWCS with TTL, and bucket partitions chronologically.
Allowing a partition to grow indefinitely (e.g. logging all events under a single
user_id) forces the JVM to allocate 200MB+ contiguous buffers during compaction, triggering multi-second Stop-The-World GC pauses and cluster gossip dropouts. Remedy: Always enforce synthetic time bucketing in the partition key.
Appending
ALLOW FILTERING to query non-indexed columns causes Cassandra to perform a brute-force sequential scan of every single SSTable across all nodes in the cluster, spiking CPU to 100%. Remedy: Create a dedicated table designed specifically around that query's exact primary key filter.
If a replica is offline for longer than
gc_grace_seconds (10 days), the rest of the cluster compacts and purges the tombstone. When the dead node restarts, its deleted data is resurrected back into the cluster. Remedy: Run nodetool repair before gc_grace_seconds elapses or decommission nodes down for >10 days.
STCS requires 50% free disk space to compact large SSTables. If your disk utilization reaches 55%, compactions fail, SSTable count explodes, and reads degrade into multi-second disk latency spirals. Remedy: Alert at 40% disk capacity or migrate to Leveled Compaction (LCS) or TWCS.