Cassandra is easy to pick for the wrong reasons. It is famous for scale, so teams choose it for a product that might one day be large, then discover that a database designed around known queries and huge write volume is painful when the queries change every sprint. It is equally easy to rule out for the wrong reasons: teams with genuinely write-heavy, always-on, multi-region workloads sometimes spend years forcing a relational primary to do something Cassandra does by design.
This article is a decision guide. It explains the architecture only as far as it changes the decision, works through a sized example, compares the realistic alternatives and lists what you sign up for when you say yes. The Postgres decision guide and the DynamoDB decision guide follow the same structure, so the three can be read side by side.
What Cassandra is, architecturally
Cassandra is a leaderless, partitioned, replicated wide-column store. Every node is equal. The partition key of each row is hashed onto a token ring, and each node owns ranges of that ring. With a replication factor of three, each partition lives on three nodes per datacenter. Any node can coordinate any request: it forwards the write to the replicas and waits for as many acknowledgements as the consistency level demands.
Inside a node, storage is a log-structured merge tree. A write appends to the commit log and inserts into a sorted in-memory memtable, and that is all. There is no read before the write, no lock and no page updated in place. Memtables flush to immutable SSTables, and background compaction merges SSTables and finally discards deleted data. The LSM tree article covers why this makes writes cheap and pushes cost to reads and compaction.
Three consequences follow, and every later section is one of them. First, writes scale nearly linearly with nodes and never contend with each other. Second, the only efficient read is one that names a partition key, because that is the only thing the ring can route. Third, there is no single leader to fail over, so a node, a rack or a whole datacenter can disappear without an outage, provided the consistency levels allow it.
The rule that decides everything: model per query
In a relational database you model entities and let the planner join them. In Cassandra you model each query as its own table. The primary key has two parts: the partition key, which decides which nodes hold the data, and the clustering columns, which decide the sort order inside a partition. A query is efficient when it supplies the full partition key and reads a contiguous slice of clustering order. Anything else is a scan across the cluster.
So the first question is not how big the data is but whether you can list your queries. If the answer is yes, and they are mostly lookups and range reads by a known key, Cassandra fits. If product managers will ask new questions of the data every month, every one of those questions is a new table, a backfill and a dual-write. The data modelling guide shows the method in detail.
Where Cassandra fits well
| Workload | Why it fits |
|---|---|
| Sustained high write volume: telemetry, events, logs, metrics | Writes are appends; capacity grows by adding nodes |
| Time series read by entity and time range | Partition by entity and time bucket, cluster by timestamp |
| Always-on services that must survive node and zone loss | Leaderless replication, no failover step |
| Active-active across regions | Each region writes locally at LOCAL_QUORUM and replicates asynchronously |
| Large key-value or wide-row lookups by known key | One routed request to a replica set |
| Data that expires | Per-row or per-table TTL, and whole-window drops with time-window compaction |
The common thread is that the access pattern is fixed and the volume or availability requirement is beyond what a single primary database can comfortably give. If neither is true, Cassandra's costs buy you nothing.
Where Cassandra strains
Ad-hoc queries and analytics are the first strain. There are no joins, aggregation is limited, and filtering on non-key columns needs an index. Storage-attached indexes in Cassandra 5.0 make secondary indexing far more practical than before, but a query that cannot be narrowed to one partition still fans out to many nodes. Teams that need reporting copy the data to an analytics system rather than querying the cluster.
Transactions are the second. Lightweight transactions give compare-and-set on a single partition using Paxos, at several times the latency of a plain write. There are no multi-partition transactions in the 5.x line. The Accord protocol, which brings multi-partition ACID transactions, is in the 6.0 line, which is pre-release as of writing; check its status before you design around it.
Delete-heavy patterns are the third. A delete writes a tombstone, which survives until compaction can safely discard it after gc_grace_seconds (864,000 seconds, ten days, by default). A queue implemented as a table, where consumers delete what they read, fills partitions with tombstones that every read must skip; by default a read warns at 1,000 tombstones and fails at 100,000. The tombstones article covers the mechanics.
Finally, small datasets. A healthy cluster needs at least three nodes per datacenter, plus repair, monitoring and people who understand compaction. For a few hundred gigabytes with modest traffic, that is all cost and no benefit.
Worked example: sizing an IoT telemetry platform
A fleet of 200,000 devices each reports one reading every 10 seconds. Readings are kept for 90 days. The product needs three queries: the latest readings for one device, a device's readings over a time range, and a live list of devices over a threshold.
Write rate: 200,000 / 10 = 20,000 rows per second, steady, around the clock. At roughly 200 bytes per row, that is 20,000 x 86,400 x 200 bytes, about 346 GB per day, and about 31 TB for 90 days. With a replication factor of three, about 93 TB on disk before compression. This is squarely a Cassandra-shaped load: relentless writes, a known key, and data that expires.
Partition size decides the key. Partitioning by device alone would grow each partition forever. Bucketing by device and day gives 8,640 readings per device per day, around 1.7 MB per partition, comfortably small. The first two queries become single-partition slices. The third, devices over a threshold right now, does not fit at all: it is a filter across every partition. The right design sends readings through a stream processor that maintains the alert list, and keeps Cassandra as the system of record.
-- Query 1: latest readings for one device, newest first.
-- Query 2: one device's readings for a time range within a day.
CREATE TABLE telemetry.readings_by_device_day (
device_id uuid,
day date, -- bucket: bounds partition size
ts timestamp,
metric text,
value double,
PRIMARY KEY ((device_id, day), ts, metric)
) WITH CLUSTERING ORDER BY (ts DESC, metric ASC)
AND compaction = {'class': 'TimeWindowCompactionStrategy',
'compaction_window_unit': 'DAYS',
'compaction_window_size': 1}
AND default_time_to_live = 7776000; -- 90 days, expired whole windows drop cleanly
-- Query 3: devices currently over threshold -> NOT this table.
-- Stream readings to an alerting service instead of scanning the cluster.The application side uses prepared statements, a token-aware policy so requests go straight to a replica, and LOCAL_QUORUM so each region acknowledges locally:
from datetime import datetime, timedelta, timezone
from cassandra.cluster import Cluster, ExecutionProfile, EXEC_PROFILE_DEFAULT
from cassandra.policies import DCAwareRoundRobinPolicy, TokenAwarePolicy
from cassandra import ConsistencyLevel
profile = ExecutionProfile(
load_balancing_policy=TokenAwarePolicy(DCAwareRoundRobinPolicy(local_dc="dc1")),
consistency_level=ConsistencyLevel.LOCAL_QUORUM,
request_timeout=2.0,
)
cluster = Cluster(["10.0.0.11", "10.0.0.12"], execution_profiles={EXEC_PROFILE_DEFAULT: profile})
session = cluster.connect("telemetry")
insert = session.prepare(
"INSERT INTO readings_by_device_day (device_id, day, ts, metric, value) VALUES (?, ?, ?, ?, ?)")
latest = session.prepare(
"SELECT ts, metric, value FROM readings_by_device_day WHERE device_id = ? AND day = ? LIMIT ?")
insert.is_idempotent = True # same key, same values: the driver may retry safely
def write_reading(device_id, ts, metric, value):
# ts is the device's own reading time, so a retry rewrites the same row
# instead of creating a second one.
return session.execute_async(insert, (device_id, ts.date(), ts, metric, value))
def last_readings(device_id, n=100):
today = datetime.now(timezone.utc).date()
rows = list(session.execute(latest, (device_id, today, n)))
if len(rows) < n: # just after midnight: top up from yesterday's bucket
rows += list(session.execute(latest, (device_id, today - timedelta(days=1), n - len(rows))))
return rowsTime-window compaction groups SSTables by day, so when a whole window expires, its files are dropped without per-row tombstone work. That only holds if rows are not updated or deleted out of order, which is why the schema is insert-only. The time-series modelling article covers bucket sizing in depth.
Consistency, in one formula
Cassandra lets each request choose how many replicas must answer. With replication factor N, write consistency W and read consistency R, a read is guaranteed to see the latest acknowledged write when R + W is greater than N. With N = 3, LOCAL_QUORUM writes and reads (2 + 2 = 4, more than 3) give read-your-writes within a datacenter and tolerate one replica down. ONE for both is faster and survives two replicas down, but may return stale data until repair or read repair catches up.
This tunability is a strength only if the team uses it deliberately. Pick consistency per query: QUORUM-style for anything a user will immediately read back, ONE for telemetry ingest where a missed sample is tolerable. Never use ALL in a request path, because a single slow replica then fails the request.
Alternatives side by side
| Need | Cassandra | Postgres | DynamoDB | ScyllaDB |
|---|---|---|---|---|
| Ad-hoc queries, joins | Poor | Excellent | Poor | Poor |
| Sustained writes beyond one primary | Excellent | Needs sharding | Excellent | Excellent |
| Multi-region active-active | Native | Hard | Global tables | Native |
| Multi-row transactions | Single partition (6.0 pre-release) | Full ACID | Limited, item count capped | Single partition |
| Operational burden | High self-hosted | Moderate | Lowest | High self-hosted |
| Cloud portability | Any cloud or on-premises | Any | AWS only | Any |
ScyllaDB speaks CQL and targets the same workloads with a different C++ implementation, so the modelling decision is the same; the choice between them is about operations, licensing and support. If you want Cassandra's data model without running it, managed options include Amazon Keyspaces and DataStax Astra DB, but check each service's compatibility notes, because not every CQL feature or consistency level behaves identically.
Failure modes in production
| Symptom | Cause | Prevention |
|---|---|---|
| One node hot, others idle | Low-cardinality partition key | Add a bucket or high-cardinality component to the key |
| Read timeouts on certain keys | Partitions grown to hundreds of MB | Bucket by time; alert on large-partition warnings |
| Tombstone warnings, then failed reads | Queue-like deletes or null inserts | Avoid delete-heavy designs; do not bind nulls |
| Deleted data reappears | Repair not run within gc_grace_seconds | Schedule repair on every table inside that window |
| Disk fills during compaction | Too little headroom for merging | Plan free space for the compaction strategy you chose |
| Cluster-wide latency spikes | ALLOW FILTERING or large IN queries | Ban both in code review; add a table instead |
What saying yes costs
You commit to query-first modelling, which means schema changes are product decisions with backfills attached. You commit to operations: repair schedules, compaction strategy choices, capacity planning with disk headroom, JVM and heap tuning, and upgrades across a rolling cluster. And you commit to a second system for analytics and for any query the tables do not serve.
The escape hatches are reasonable. Data leaves Cassandra easily through bulk export or Spark, and the same CQL tables move to ScyllaDB or a managed CQL service without remodelling. What does not transfer is the access-pattern discipline, and that is the real asset: teams that do it well end up with a write path that does not need re-architecting as traffic grows by orders of magnitude.
What to do next
- Write down every query the product needs, with its expected rate and latency target. If you cannot, choose Postgres now and revisit later.
- Estimate the write rate and retained volume with the arithmetic above. If it fits comfortably on one primary with replicas, Cassandra is probably premature.
- For each query, draft a table with a partition key and clustering order, and compute the worst-case partition size.
- Identify queries that need cross-partition filters or transactions, and decide which other system will serve them.
- Choose consistency levels per query and confirm that R + W exceeds the replication factor wherever users read their own writes.
- Before production, schedule repair, alert on tombstone and large-partition warnings, and load-test at twice the expected peak with one node down.