Cassandra: Wide-Column at Scale

Your clickstream pipeline is humming: Kafka carries the events, your streaming job counts them, and product just asked for the obvious next feature — "recent activity" on the profile page. The query is simple: give me the last 50 events for user X, served in milliseconds, while the write firehose never stops. Your Postgres can answer it today with an index on (user_id, event_time) — and at ten times the write volume, with a second data center in the plan, you start pricing out what that index costs you on every single insert.

This post introduces the store built for exactly that shape of work: Apache Cassandra, the wide-column database that trades ad-hoc querying for write throughput and scale. But Cassandra's real lesson isn't a new query language — it's a new design habit. In a relational database you model the entities and query however you like afterward. In Cassandra you do the opposite: you model the query first, and the table follows. Get that backwards and Cassandra will simply refuse your queries — this post shows you the refusal, verbatim, and then the table design that earns the query back.

Lab honesty, up front. Everything CQL below ran for real: Apache Cassandra 5.0.9, single node, Temurin JDK 17, cqlsh 6.2.0 — every command shown was executed and every output pasted verbatim (long DESCRIBE output is trimmed with an explicit marker). Two things did not run in this lab, and the post says so where it matters: nodetool status couldn't run because this sandbox blocks raw TCP for Java processes (a sandbox policy, not a Cassandra problem — the node itself is healthy and serving CQL), and the Java driver program compiled cleanly with javac but couldn't connect for the same reason. The program is complete and the exact run command is given; run it on your own machine.

CQL looks like SQL (and isn't)

Cassandra's query language, CQL, is deliberately SQL-flavored — and deliberately not SQL. No joins, no subqueries, no GROUP BY across partitions, no ad-hoc WHERE on any column you fancy. The resemblance is a kindness to newcomers; the restrictions are the whole point. They exist because every CQL query must be answerable by contacting a small, known set of nodes — and that constraint is what lets Cassandra scale writes linearly.

First, a keyspace — Cassandra's equivalent of a database, and the place where you declare the replication strategy:

CREATE KEYSPACE IF NOT EXISTS user_events
  WITH replication = {'class': 'SimpleStrategy', 'replication_factor': 1};

DESCRIBE KEYSPACE user_events;
CREATE KEYSPACE user_events WITH replication = {'class': 'SimpleStrategy', 'replication_factor': '1'}  AND durable_writes = true;

Two things to notice. SimpleStrategy with replication factor 1 means "keep one copy, don't think about racks or data centers" — the right choice for a single-node lab and only for a single-node lab. In production, especially across data centers, you would use NetworkTopologyStrategy and say how many copies belong in each data center; SimpleStrategy is unaware of racks and DCs and will happily place all your replicas in one of them.

Now the table. Your relational instinct — the one this post is about to retrain — says: model the entity. An event has an id, a user, a type, a time, a payload. So:

USE user_events;
CREATE TABLE events (
  event_id uuid PRIMARY KEY,
      user_id text,
      event_type text,
  event_time timestamp,
  payload text
);
INSERT INTO events (event_id, user_id, event_type, event_time, payload) VALUES (11111111-1111-1111-1111-111111111111, 'u1', 'click', '2026-10-01 10:00:00', 'product 42');
INSERT INTO events (event_id, user_id, event_type, event_time, payload) VALUES (22222222-2222-2222-2222-222222222222, 'u1', 'purchase', '2026-10-01 10:05:00', 'order 7');
INSERT INTO events (event_id, user_id, event_type, event_time, payload) VALUES (33333333-3333-3333-3333-333333333333, 'u2', 'click', '2026-10-01 10:01:00', 'product 9');

Three rows, tiny table, and now the product query — last 50 events for user X:

SELECT user_id, event_type, event_time FROM events WHERE user_id = 'u1';
InvalidRequest: Error from server: code=2200 [Invalid query] message="Cannot execute this query as it might involve data filtering and thus may have unpredictable performance. If you want to execute this query despite the performance unpredictability, use ALLOW FILTERING"

Read that error as Cassandra teaching you its data model. With event_id as the primary key, each row lives on whichever node owns the hash of its event_id — u1's events are scattered across the cluster. To find "all events for u1," Cassandra would have to ask every node to scan all its data. On three rows that's instant; on three billion it's a full-cluster scan that can time out and drag down every other query while it runs. So Cassandra refuses — unless you insist:

SELECT user_id, event_type, event_time FROM events WHERE user_id = 'u1' ALLOW FILTERING;
 user_id | event_type | event_time
---------+------------+---------------------------------
      u1 |   purchase | 2026-10-01 10:05:00.000000+0000
      u1 |      click | 2026-10-01 10:00:00.000000+0000

(2 rows)

It works — on three rows. Treat ALLOW FILTERING as a rule of thumb, not a tool: fine for exploring a small table in cqlsh, a design smell in application code. If your application needs the query, the table is wrong. Fix the table, not the query.

The main event: model the query first

Here is the habit to build. Before writing CREATE TABLE, write down the queries the table must serve — all of them, including the ones product will ask for next quarter. Our list has exactly one entry:

Q1: the last 50 events for a given user, newest first.

Now design the primary key from the query outward. A Cassandra primary key has two jobs, and they map directly onto the two halves of the query:

  • The partition key — the first part of the primary key — decides which node holds the data and what you can look up by. Cassandra hashes it (Murmur3 by default) to a token, and the token decides the node. In the primary-key access pattern, every query names the full partition key, because that's how Cassandra knows which handful of nodes to ask instead of all of them. Our query looks up by user — so the partition key is user_id. One user's events, one partition, one (or a few, with replicas) nodes.
  • The clustering columns — the rest of the primary key — decide the sort order of rows inside the partition. Our query wants newest first — so cluster by event_time descending. Two events can share a millisecond, so add event_id as a tiebreaker to keep the key unique.
Query-first: the query designs the table Q1: last 50 events for user X, newest first write the query before the CREATE TABLE PRIMARY KEY (user_id, event_time, event_id) the table follows the query Partition key: user_id Murmur3(user_id) → token → node u1's events live together on the nodes owning that token range ✓ WHERE user_id = 'u1' hits one partition ✗ WHERE event_type = 'click' has no partition key → refused Clustering: event_time DESC rows sorted on disk, newest first 10:05 purchase → 10:03 click → 10:00 click event_id ASC breaks same-ms ties ✓ "last 50" = read the first 50 rows ✓ time ranges slice the sorted rows

The table that falls out of Q1:

CREATE TABLE events_by_user (
  user_id text,
  event_time timestamp,
  event_id uuid,
  event_type text,
  payload text,
  PRIMARY KEY (user_id, event_time, event_id)
) WITH CLUSTERING ORDER BY (event_time DESC, event_id ASC);
CREATE TABLE user_events.events_by_user (
    user_id text,
    event_time timestamp,
    event_id uuid,
    event_type text,
    payload text,
    PRIMARY KEY (user_id, event_time, event_id)
) WITH CLUSTERING ORDER BY (event_time DESC, event_id ASC)
    [... table options elided: compaction, caching, compression, gc_grace_seconds, ...]
    AND speculative_retry = '99p';

Load it — deliberately inserting out of order — and ask the product query:

INSERT INTO events_by_user (user_id, event_time, event_id, event_type, payload) VALUES ('u1', '2026-10-01 10:00:00', 11111111-1111-1111-1111-111111111111, 'click', 'product 42');
INSERT INTO events_by_user (user_id, event_time, event_id, event_type, payload) VALUES ('u1', '2026-10-01 10:05:00', 22222222-2222-2222-2222-222222222222, 'purchase', 'order 7');
INSERT INTO events_by_user (user_id, event_time, event_id, event_type, payload) VALUES ('u1', '2026-10-01 10:03:00', 44444444-4444-4444-4444-444444444444, 'click', 'product 13');

SELECT user_id, event_time, event_type, payload FROM events_by_user WHERE user_id = 'u1' LIMIT 50;
 user_id | event_time                      | event_type | payload
---------+---------------------------------+------------+------------
      u1 | 2026-10-01 10:05:00.000000+0000 |   purchase |    order 7
      u1 | 2026-10-01 10:03:00.000000+0000 |      click | product 13
      u1 | 2026-10-01 10:00:00.000000+0000 |      click | product 42

(3 rows)

Inserted 10:00, 10:05, 10:03 — returned 10:05, 10:03, 10:00. The clustering order holds regardless of insert order, because rows are stored sorted on disk. "Last 50" is now literally "read the first 50 rows of the partition" — the cheapest possible read. And time ranges slice the sorted rows for free:

SELECT user_id, event_time, event_type FROM events_by_user
  WHERE user_id = 'u1' AND event_time > '2026-10-01 10:02:00' AND event_time < '2026-10-01 10:06:00';
 user_id | event_time                      | event_type
---------+---------------------------------+------------
      u1 | 2026-10-01 10:05:00.000000+0000 |   purchase
      u1 | 2026-10-01 10:03:00.000000+0000 |      click

(2 rows)

But stray from the designed query and the refusal is back — same error as before, now on the "right" table:

SELECT user_id, event_time FROM events_by_user WHERE event_time > '2026-10-01 10:02:00';
SELECT user_id, event_time FROM events_by_user WHERE event_type = 'click';
InvalidRequest: Error from server: code=2200 [Invalid query] message="Cannot execute this query as it might involve data filtering and thus may have unpredictable performance. If you want to execute this query despite the performance unpredictability, use ALLOW FILTERING"
InvalidRequest: Error from server: code=2200 [Invalid query] message="Cannot execute this query as it might involve data filtering and thus may have unpredictable performance. If you want to execute this query despite the performance unpredictability, use ALLOW FILTERING"

This is the discipline, stated plainly: a Cassandra table serves the queries it was designed for; anything else is refused outright or reduced to a cluster-wide scan. Need "all clicks across users"? That's a different query — so it's a different table (say, events_by_type partitioned by event type), with the same events duplicated into it. Duplicating data per query feels wrong to a normalized-schema brain; in Cassandra it's the standard practice, a rule of thumb you'll see everywhere: storage is cheap, joins don't exist, so you denormalize at write time — one table per query. Writes are cheap and idempotent, which is what makes maintaining two or three copies affordable.

One more design pressure to respect: keep partitions bounded. The partition key decides what lives together, and a partition that's too wide — one user with ten years of events, or worse, a partition key like "country" holding half your data — creates hot spots and huge single partitions that are expensive to read and repair. A common rule of thumb is to keep partitions in the tens-of-megabytes range and add a time bucket (say, month) to the partition key when a single key would grow without bound. Rules of thumb, not laws — measure your own partitions.

Principle: model the query first; the table follows. Write the queries before the CREATE TABLE. If you can't name the partition key from the query's WHERE clause, you don't have a design yet.

Tunable consistency: a dial, not a doctrine

You'll hear Cassandra summarized as "the AP database." Treat that as architectural shorthand, explicitly labeled — never as the whole consistency model. What Cassandra actually gives you is a dial you turn per operation: every read and write carries a consistency level naming how many replicas must acknowledge. The cluster's bias is toward staying available under partitions, but the consistency behavior of your query is yours to choose, each time.

Watch the dial turn on our single node (replication factor 1, so levels above ONE have nowhere to go — which is exactly what makes the failure below instructive):

CONSISTENCY;
CONSISTENCY QUORUM;
SELECT count(*) FROM events_by_user;
Current consistency level is ONE.
Consistency level set to QUORUM.

 count
-------
     5

(1 rows)

Warnings :
Aggregation query used without partition key

(The warning is cqlsh telling you count(*) scanned partitions — honest on one node, expensive on a hundred. It's a warning, not an error.)

Writes take a level too. ALL on a one-replica keyspace just means "the one replica":

CONSISTENCY ALL;
INSERT INTO events_by_user (user_id, event_time, event_id, event_type, payload) VALUES ('u3', '2026-10-01 11:00:00', 66666666-6666-6666-6666-666666666666, 'signup', '-');
SELECT user_id FROM events_by_user WHERE user_id = 'u3';
Consistency level set to ALL.

 user_id
---------
      u3

(1 rows)

Now ask for more replicas than exist:

CONSISTENCY TWO;
SELECT count(*) FROM events_by_user;
NoHostAvailable: ('Unable to complete the operation against any hosts', {<Host: 127.0.0.1:9042 datacenter1>: Unavailable('Error from server: code=1000 [Unavailable exception] message="Cannot achieve consistency level TWO" info={'consistency': 'TWO', 'required_replicas': 2, 'alive_replicas': 1}')})

required_replicas: 2, alive_replicas: 1 — the whole model in one error. A consistency level is a count of replicas that must respond; if fewer are alive, the operation fails loudly rather than answering weakly. That's the dial:

  • ONE — one replica answers. Fastest, least coordination. During a partition or a downed node, the price is that you may read stale data or your write may sit on one replica until the others catch up (via hints and repair). The usual pick for high-throughput, low-stakes data — counters, session state, event streams.
  • QUORUM — a majority of replicas (2 of 3, 3 of 5). The workhorse default, as a rule of thumb: it keeps working through a single replica failure while making stale reads much less likely.
  • ALL — every replica. Strongest guarantee, and the operation fails if any replica is unreachable. You pay in availability for certainty — the right call for the occasional operation where correctness outweighs uptime, not for the hot path.

The classic rule of thumb for "read-your-writes" style guarantees: write consistency level + read consistency level > replication factor (e.g. write at QUORUM, read at QUORUM, RF 3). It's a useful majority-overlap argument, not a law — it assumes the same replica set and healthy repair, and real guarantees come from measuring, not arithmetic.

And the AP shorthand, stated carefully: with RF 3 across racks, if a network partition isolates one replica, writes at ONE or QUORUM keep succeeding on the reachable side while the isolated replica catches up later through hinted handoff and repair. The architecture leans toward availability — it won't refuse writes just because it can't reach everyone — but nothing stops you from demanding ALL and failing instead. The bias is in the defaults and the machinery; the decision is per query.

Under the hood: the LSM-tree (one honest paragraph)

This section is conceptual — standard, documented storage-engine design, not something this lab measured. Cassandra stores data in an LSM-tree: writes land first in an in-memory memtable (and are appended to the commit log for durability), and when the memtable fills it is flushed to disk as an immutable, sorted SSTable. Nothing is ever updated in place — a later write to the same key is just a newer entry in a newer SSTable, with tombstones marking deletes. A background process called compaction periodically merges SSTables, discarding overwritten values and tombstones past their grace period. That's why Cassandra writes so fast (every write is a sequential append) and why reads sometimes check several SSTables plus in-memory bloom filters before answering. You saw a hint of this machinery in the DESCRIBE TABLE output earlier: compaction = {'class': '...SizeTieredCompactionStrategy', ...} — the default merge policy for the table.

The LSM-tree write path (conceptual) WRITE append, never update in place Memtable in RAM, sorted + commit log on disk SSTables immutable, sorted files Compaction merges, drops old versions flush Reads merge the memtable, relevant SSTables, and bloom filters — newer timestamps win.

Cassandra from Java: the DataStax driver

CQL is for exploration; applications talk to Cassandra through a driver. The standard one is the DataStax Java driver 4.x — on Maven Central as com.datastax.oss:java-driver-core:4.17.0 (this lab used 4.17.0; note the core jar is not shaded, so its runtime dependencies — Netty, Typesafe config, Jackson, Dropwizard Metrics, and friends — ride along on the classpath). The driver speaks Cassandra's native binary protocol, manages a connection pool per node, and handles token-aware routing: it learns which nodes own which token ranges and sends each query straight to a replica when it can.

Here is the complete program — connect, insert, prepared-statement select, delete, close:

import com.datastax.oss.driver.api.core.CqlSession;
import com.datastax.oss.driver.api.core.cql.BoundStatement;
import com.datastax.oss.driver.api.core.cql.PreparedStatement;
import com.datastax.oss.driver.api.core.cql.ResultSet;
import com.datastax.oss.driver.api.core.cql.Row;
import com.datastax.oss.driver.api.core.metadata.Node;

import java.net.InetSocketAddress;
import java.util.UUID;

/**
 * Minimal CRUD against the events_by_user table from the post-7 lab.
 * Compile: javac -cp "lib/*" CassandraCrud.java
 * Run:     java -cp "lib/*:." CassandraCrud
 */
public class CassandraCrud {

    public static void main(String[] args) {
        try (CqlSession session = CqlSession.builder()
                .addContactPoint(new InetSocketAddress("127.0.0.1", 9042))
                .withLocalDatacenter("datacenter1")
                .build()) {

            System.out.println("Connected. Cluster nodes:");
            for (Node node : session.getMetadata().getNodes().values()) {
                System.out.println("  " + node.getEndPoint()
                        + " dc=" + node.getDatacenter()
                        + " cassandra=" + node.getCassandraVersion()
                        + " state=" + node.getState());
            }

            session.execute("USE user_events");

            // CREATE
            UUID id = UUID.fromString("77777777-7777-7777-7777-777777777777");
            session.execute(
                    "INSERT INTO events_by_user (user_id, event_time, event_id, event_type, payload) " +
                    "VALUES ('u9', '2026-10-06 12:00:00', " + id + ", 'login', 'driver demo')");
            System.out.println("INSERT ok: u9 login event @ 2026-10-06 12:00:00");

            // READ (prepared + bound statement, the way production code does it)
            PreparedStatement ps = session.prepare(
                    "SELECT user_id, event_time, event_type, payload FROM events_by_user " +
                    "WHERE user_id = ? LIMIT 50");
            BoundStatement bound = ps.bind("u9");
            ResultSet rs = session.execute(bound);
            System.out.println("SELECT for u9:");
            int n = 0;
            for (Row row : rs) {
                n++;
                System.out.println("  " + row.getString("user_id")
                        + " | " + row.getInstant("event_time")
                        + " | " + row.getString("event_type")
                        + " | " + row.getString("payload"));
            }
            System.out.println("  (" + n + " row(s))");

            // DELETE + verify
            session.execute("DELETE FROM events_by_user WHERE user_id = 'u9' "
                    + "AND event_time = '2026-10-06 12:00:00' AND event_id = " + id);
            long remaining = session.execute(
                    "SELECT count(*) FROM events_by_user WHERE user_id = 'u9'")
                    .one().getLong(0);
            System.out.println("DELETE ok. Rows left for u9: " + remaining);
        }
        System.out.println("Session closed.");
    }
}

Three things worth noticing before you run it. withLocalDatacenter("datacenter1") names the local DC so the driver can prefer nearby replicas — datacenter1 is Cassandra's default DC name, the one this lab's node reports. try (CqlSession ...) matters: the session owns Netty event loops and connection pools, and closing it is what releases them. And the read uses a prepared statement — the CQL is parsed once server-side and re-executed with bound values, which is both faster and the standard defense against CQL injection, the same way PreparedStatement works in JDBC.

Compile (real — javac exited 0 with no errors, producing CassandraCrud.class):

$ javac -cp "lib/*" CassandraCrud.java
# exit 0, no errors — CassandraCrud.class produced

[NOT RUN] — sandbox note, honest: this lab's sandbox blocks raw TCP connections from Java processes (a runtime policy — the message names the "Direct network protocols" permission), so the driver can't reach port 9042 here even though cqlsh and the node are fine. The run step is yours, on any machine with the node from this post running:

$ java -cp "lib/*:." CassandraCrud

Reading the code, that's the full CRUD loop against the query-first table: connect, print node metadata, insert, prepared-statement read, delete, verify zero rows remain, close the session.

Cheat sheet

  • Model the query first; the table follows. Write the queries before CREATE TABLE; if the partition key isn't visible in the WHERE clause, you don't have a design yet.
  • Partition key = distribution + lookup key (hashed to a token → node); in the primary-key access pattern, every query names it. Clustering columns = sort order within the partition, stored sorted on disk.
  • CQL looks like SQL and isn't: no joins, no ad-hoc WHERE. Queries that would scan the cluster are refused — ALLOW FILTERING is for cqlsh exploration, a design smell in application code.
  • One table per query is the standard practice: denormalize at write time, since storage is cheap and there are no joins to lean on. Keep partitions bounded (rule of thumb: tens of MB; add a time bucket to the partition key when one key would grow forever).
  • Consistency is a per-operation dial: ONE (fastest), QUORUM (the workhorse default, as a rule of thumb), ALL (strongest, fails if any replica is down). "Leans AP" is architectural shorthand — the behavior is yours to choose each time.
  • Write CL + read CL > RF is a useful rule of thumb for read-your-writes, not a law.
  • LSM-tree: memtable → immutable sorted SSTables → compaction merges; writes are sequential appends, which is where the write throughput comes from.
  • SimpleStrategy is single-node-only; production clusters use NetworkTopologyStrategy with per-DC replication factors.
  • DataStax java-driver 4.x: CqlSession.builder().addContactPoint(...).withLocalDatacenter(...), prepared statements for repeated queries, close the session (try-with-resources).
  • Cassandra 5.0.x wants JDK 11/17 — match the tool's runtime, not the newest JDK.

What's next

You've now felt Cassandra's central trade: it will serve a designed query at enormous scale and simply refuse an undesigned one. That refusal is a feature — it's the database telling you where your model is incomplete. But "tunable consistency" is still consistency you chose per operation, and the machinery underneath (hints, repair, read-repair) is doing real work to paper over the gaps. What if you want the database to guarantee strong per-row consistency instead of asking you to tune it?

The next post, HBase: Strong Consistency on Hadoop, is Cassandra's mirror image: same wide-column shape, opposite philosophy — one region server owns each row range, coordinated infrastructure underneath, strong row-level consistency by construction. You'll design row keys the way you just designed partition keys, and learn exactly what the coordination costs.

Field check before you move on: (1) Without running it, predict which of these get refused and why: WHERE user_id='u1' AND event_type='click', WHERE user_id IN ('u1','u2'), WHERE event_time > '2026-10-01 10:00:00' ALLOW FILTERING — then run them and check. (2) Change the table's CLUSTERING ORDER to event_time ASC, re-run the "last 50" query — what does your application have to do now, and what did you lose? (3) Design a second table for the query "all purchase events across users in the last hour" — what's its partition key, and why is event_type alone a dangerous one? (4) On your own machine, run the Java program above and add one line: print the consistency level used by the prepared SELECT (hint: check the driver's Statement.setConsistencyLevel), then re-run it at ConsistencyLevel.ALL.

Continue: Java Learning Roadmap 2026

Comments

Popular posts from this blog

JSP Servlet Interview Questions For Freshers Series 1

Java Banking Finance Services and Insurance (BFSI) domain interview questions

Java program to check even or odd number