Apache Cassandra is a free, open-source, distributed NoSQL database of the wide-column type, built to store very large volumes of data across many servers with no single point of failure. It was developed at Facebook by Avinash Lakshman and Prashant Malik to power inbox search, released as open source in July 2008, and has been a top-level Apache Software Foundation project since February 2010. The current major version is Cassandra 5.0. A Cassandra DB trades the joins and ad-hoc queries of a relational database for fast writes, linear scaling and the ability to keep running when whole servers or data centres go down.
How the Cassandra database works
A masterless ring
Every node in a Cassandra cluster has the same role. There is no primary server that others depend on, so any node can accept any read or write; the node that receives a request acts as the coordinator and forwards it to the nodes that own the data. Nodes learn about each other’s health through a gossip protocol. Losing one node, or even a whole rack, does not stop the cluster.
Partitioner and token ring
Each row belongs to a partition, identified by its partition key. The partitioner (Murmur3Partitioner by default) hashes the partition key to a token, a 64-bit number. The full token range is drawn as a ring, and each node owns slices of it. A row goes to the node whose slice contains its token. Because the hash spreads keys evenly, adding a node takes a share of the ring and of the data from the others, which is why capacity grows roughly in step with node count.
Replication factor
The replication factor (RF) is how many copies of each partition the cluster keeps. With RF = 3, the partition is stored on the node that owns its token and on the next two replicas chosen by the replication strategy. NetworkTopologyStrategy lets you set RF per data centre, for example 3 in Mumbai and 3 in Singapore.
Tunable consistency, with a QUORUM example
For each query the application chooses a consistency level: how many replicas must reply before the coordinator reports success. Common levels are ONE, QUORUM, LOCAL_QUORUM and ALL. A quorum is a majority of replicas:
QUORUM = floor(RF / 2) + 1
Worked example with RF = 3:
- QUORUM = floor(3/2) + 1 = 1 + 1 = 2 replicas.
- Write at QUORUM (W = 2) and read at QUORUM (R = 2). Then R + W = 4 > RF = 3, so the set of replicas that acknowledged the write and the set that answered the read must share at least one node. The read therefore sees the latest acknowledged write.
- Both operations still succeed with one of the three replicas down, because 2 are still available.
- Compare ONE/ONE: R + W = 1 + 1 = 2, which is not greater than 3. Reads are faster but can return stale data until the replicas catch up. This is called eventual consistency.
With RF = 5, QUORUM = floor(5/2) + 1 = 3, and the cluster tolerates two replicas down for QUORUM operations.
The Cassandra data model
- Keyspace: the top-level container, similar to a database or schema. The replication settings are defined here.
- Table: rows and columns, defined in CQL (Cassandra Query Language).
- Partition key: decides which nodes store the row. All rows with the same partition key live together.
- Clustering columns: decide the sort order of rows inside a partition, so range reads within a partition are fast.
CREATE TABLE sensor_readings (
sensor_id text,
reading_ts timestamp,
temp_c double,
PRIMARY KEY ((sensor_id), reading_ts)
) WITH CLUSTERING ORDER BY (reading_ts DESC);Here sensor_id is the partition key and reading_ts the clustering column. “Latest 100 readings for sensor S-17” is one fast read from one partition. “All sensors above 40°C” is not something this table can answer efficiently.
That is the core rule of Cassandra: query-first modelling. You list the queries the application will run and design one table per query, copying data between tables where needed. Denormalisation and duplication are normal, the opposite of relational practice. A partition should also stay bounded in size (a common guideline is under about 100 MB), so time-series data is often bucketed, for example by sensor and day.
Is Cassandra a relational database?
No. Cassandra is a NoSQL database, even though CQL looks a lot like SQL (CREATE TABLE, SELECT, INSERT, WHERE). The differences matter in practice:
- No joins and no foreign keys. Related data is duplicated into the tables that need it.
- The WHERE clause is restricted. You filter on the partition key and then on clustering columns in order. Filtering on other columns needs an index or the
ALLOW FILTERINGoption, which can scan the cluster and is usually a design mistake. - No multi-row ACID transactions in the relational sense. Writes to one partition are atomic and isolated; lightweight transactions (
IF NOT EXISTS,IF col = value) give compare-and-set on a single partition at extra cost.
Put simply, a relational database lets you ask new questions of existing data. Cassandra answers the questions you designed for, at very large scale.
Pros of Cassandra, and why
- Very fast writes. The write path is a log-structured merge (LSM) design. A write is appended to the commit log on disk for durability and written to the memtable in memory, and that is it. Memtables are later flushed to immutable SSTables. There is no read-before-write and no in-place update on disk.
- Linear scalability. Add nodes and the token ring rebalances. Read and write throughput rise roughly in proportion.
- High availability. No master, data replicated RF times, and consistency you can relax to stay up during failures.
- Multi-data-centre replication built in, useful for low latency across regions and for disaster recovery.
- Flexible schema within a table. Columns can be added without rewriting existing rows, and collection types (list, set, map) are supported.
Cons of Cassandra, and why
- No ad-hoc queries. A new query pattern usually means a new table and a data backfill.
- Reads cost more than writes. A row may be spread across the memtable and several SSTables, which must be merged at read time. Bloom filters and caches help, but reads are the weaker side.
- Compaction. Background merging of SSTables uses disk I/O and needs spare disk space. A badly chosen compaction strategy can hurt latency.
- Tombstones. A delete writes a marker (a tombstone) instead of removing data. Tombstones stay until compaction after
gc_grace_seconds(10 days by default), and a read that scans many of them slows down or fails. Heavy delete or queue-like workloads suit Cassandra badly. - Eventual consistency by default. Strong consistency is possible with QUORUM or ALL, but it costs latency and availability.
Cassandra use cases
| Use case | Why Cassandra fits |
|---|---|
| IoT and sensor time series | Huge write volume, reads by device and time range |
| Messaging and chat history | Messages partitioned by conversation, sorted by time |
| User activity, event logs and feeds | Append-heavy, read by user |
| Product catalogues and personalisation | High read and write traffic, always-on across regions |
| Fraud and security event stores | Fast ingest, lookups by account or device |
| Not a fit: banking ledgers, reporting, ad-hoc analytics | Need joins, multi-row transactions or flexible queries |
The Apache project lists users such as Apple, Netflix, Uber, Spotify, eBay, Walmart and the Indian fantasy sports platform Dream11 among its case studies.
Cassandra vs MongoDB vs a relational database
| Feature | Cassandra | MongoDB | Relational (MySQL, PostgreSQL) |
|---|---|---|---|
| Data model | Wide-column (partitions of rows) | JSON-like documents | Normalised tables |
| Architecture | Masterless, every node equal | One primary per replica set, secondaries follow | Usually one primary, read replicas |
| Query flexibility | Low, designed per query | High, rich queries and indexes | Highest, full SQL with joins |
| Transactions | Single-partition only, plus lightweight transactions | Multi-document ACID supported | Full ACID |
| Best at | Massive writes, multi-region uptime | Flexible app data, fast development | Consistent business data, reporting |
Cassandra is often used next to Hadoop or Spark: Cassandra serves live traffic while the batch system does the heavy analysis. The trade-offs on that side are covered in Hadoop pros and cons, and the broader families of database systems in classifications of DBMS.
References
- Apache Cassandra, official documentation, downloads and case studies.
- SQL: From Traditional Databases to Big Data, ResearchGate.
FAQs
What is Cassandra DB used for?
Cassandra is used for applications that write large volumes of data and must stay available at all times: IoT and sensor time series, messaging history, activity logs, feeds, catalogues and fraud event stores. It is a poor fit for workloads that need joins, multi-row transactions or ad-hoc reporting.
Is Cassandra SQL or NoSQL?
Cassandra is a NoSQL wide-column database. Its query language, CQL, uses SQL-like keywords, but it has no joins, no foreign keys and restricted WHERE clauses, so tables must be designed around the queries the application will run.
What is QUORUM in Cassandra?
QUORUM is a consistency level requiring a majority of replicas, floor(RF/2) + 1, to respond. With a replication factor of 3, QUORUM is 2. Writing and reading at QUORUM gives R + W = 4, greater than 3, so reads return the latest acknowledged write while one replica can be down.
Who created Cassandra?
Avinash Lakshman and Prashant Malik built Cassandra at Facebook for inbox search. Facebook released it as open source in July 2008, it entered the Apache Incubator in March 2009 and became a top-level Apache project in February 2010.
Why are deletes a problem in Cassandra?
A delete writes a tombstone marker instead of removing data. Tombstones remain until they are older than gc_grace_seconds (10 days by default) and compaction removes them, and reads that scan many tombstones become slow or fail. Workloads with heavy deletes need careful design.
