A Cassandra materialized view is a second table that the server keeps in step with a base table, keyed differently so you can query the same rows by another column. It looks like the answer to Cassandra's most common modelling problem, needing two access paths to one entity. It is also a feature that the project's own configuration file describes as experimental and not recommended for production use, and that ships disabled.

Both statements are true, and the reason is the write path. This article walks through exactly what happens when you write to a table with a view, why the view can silently fall out of step with its base, how to detect that, and when an application-maintained table or a storage-attached index is the better choice.

Advertisement

What a view is, and the rules it must follow

A view is defined by a SELECT over one base table with a new primary key. Cassandra imposes two rules on that key. It must contain every primary key column of the base table, so that each view row corresponds to exactly one base row. And it may add at most one column that is not part of the base primary key. The WHERE clause must also restrict every view primary key column with IS NOT NULL, because a row cannot be stored under a null key.

CREATE TABLE app.users (
    user_id  uuid PRIMARY KEY,
    email    text,
    name     text,
    country  text,
    created  timestamp
);

-- View PK = the base PK (user_id) + at most ONE non-PK base column (email).
CREATE MATERIALIZED VIEW app.users_by_email AS
    SELECT user_id, email, name
    FROM app.users
    WHERE email IS NOT NULL AND user_id IS NOT NULL
    PRIMARY KEY ((email), user_id);

-- Reads by the new key
SELECT user_id, name FROM app.users_by_email WHERE email = 'ana@example.com';

-- Operations
-- nodetool viewbuildstatus app users_by_email
-- SELECT * FROM system_distributed.view_build_status WHERE keyspace_name = 'app';

The view is a real table: it has its own SSTables, its own token for each row, and its own replicas, which are usually different nodes from the base row's replicas. That last fact drives everything else.

The write path, step by step

CoordinatorUPDATE users ... CL=QUORUMBase replica A (1st)lock, read, diffBase replica B (2nd)lock, read, diffBase replica C (3rd)lock, read, diffView replica X (1st)for view tokenView replica Y (2nd)for view tokenView replica Z (3rd)for view tokenbase mutationvia local batchlogvia local batchlogvia local batchlogview delta: tombstone old view row + insert new onebase replicas ack the coordinator without waiting for the view replicas
Writing to a base table with a view: each base replica computes the view change locally and sends it to the view replica with the same position in its replica list, through its local batchlog.

A client write to the base table goes to a coordinator, which sends it to the base replicas as usual; see the write and read path for that part. What changes is on each base replica:

  1. The replica takes a lock on the base partition, so that concurrent updates to the same row compute view changes in a consistent order.
  2. It reads the current local version of the row. This read-before-write is what makes a view write far more expensive than an ordinary Cassandra write, which never reads.
  3. It computes the view delta. If the change affects the view key, for example a new email, the delta is a tombstone for the old view row and an insert for the new one. If it only changes a selected column, the delta is an update to the existing view row.
  4. It works out its paired view replica. Cassandra pairs by position: the replica that is first in the base row's replica list, restricted to the local datacenter, sends to the first replica in the view row's list, the second to the second, and so on. Nodes that replicate both the base and the view token are excluded from the pairing and apply the view change locally.
  5. If the paired replica is this node, the view mutation is applied locally. Otherwise it is sent to the paired replica through the local batchlog, so that if the send fails it can be replayed later. When no pairing exists, for example during a range movement, the mutation goes to the local batchlog for replay.
  6. The base mutation is applied and acknowledged to the coordinator. The base replica does not wait for the view replica to acknowledge.

The result is that the client's consistency level describes the base write only. A QUORUM write that succeeds guarantees the base row on a quorum; the view rows follow asynchronously, one pair at a time, and each base replica reads its own local version of the row, which may itself be stale.

Advertisement

Reading from a view

You read a view like any table, but you cannot write to it; all changes come through the base. Reads use ordinary consistency levels, and those levels apply to the view's replicas only. A QUORUM read of the view straight after a QUORUM write to the base is not guaranteed to see that write, because the view update may still be in flight or in a batchlog waiting for replay. Read-your-writes, which careful Cassandra users get from QUORUM on both sides, does not hold across base and view.

Design readers accordingly. A login flow that looks users up by email should tolerate a missing or outdated row, for example by confirming the result against the base table by primary key before trusting it. That second read is cheap, because it goes to a single partition, and it turns a silent wrong answer into a detectable miss. If every caller needs that guard, it is a strong sign the view is the wrong tool for the job.

A worked example: changing an email

User 7 changes their email from ana@old.example to ana@new.example. On each base replica the lock is taken, the old row is read (email ana@old.example), and the delta is computed: delete the view row in partition ana@old.example and insert one in partition ana@new.example. Those two view partitions have different tokens, so each may live on a different set of view replicas, and the pairing is computed per view token.

Now suppose one base replica missed an earlier update and still holds an older email, ana@older.example. It computes a delta that deletes a view row under ana@older.example, which does nothing useful, and inserts the new one. The other replicas delete the correct old row. If the replicas that would have deleted ana@old.example on some view replica were the ones that failed, the view keeps a stale row pointing email ana@old.example at user 7. A lookup by the old email now returns a user who no longer has it. Nothing in the normal read path will notice.

This is the heart of the problem: the view is derived from each base replica's local view of the data, not from the agreed state, and there is no built-in process that compares base and view.

Why base and view diverge

  • Stale base replicas: a base replica that missed writes computes wrong deltas, as in the example. Run base-table repair regularly so replicas agree before they compute; see repair.
  • Repair does not compare base with view: repairing the view table reconciles view replicas with each other. If every view replica is missing a row, or every one has a stale row, repair has nothing to fix.
  • Asynchronous view writes: between the base acknowledgement and the view write, readers of the view see old data. Under load or with a slow view replica that gap grows, and batchlog replay after failures can make it longer.
  • Unselected column deletions: the documentation warns that deleting or nulling a base column that is not selected in the view may shadow missed updates to other columns, and advises against doing it.
  • Streaming shortcuts: in the 5.0 configuration, materialized_views_on_repair_enabled (default true) sends streamed base data through the write path so views are updated. Turning it off makes streaming faster but, as the configuration comment warns, in extreme cases data can reach the base SSTables and never the view.

Building a view on existing data

Creating a view on a populated table starts a view build: each node scans its local base data and generates view mutations. It runs with concurrent_materialized_view_builders threads, 1 by default, and can take hours on a large table while adding read and write load. Track it with nodetool viewbuildstatus or the system_distributed.view_build_status table, and do not send reads to the view until every node reports it as built. If a build is interrupted it resumes, but plan view creation like a data migration, with capacity headroom and a maintenance window.

What a view costs

Each write to the base now includes a local read, a partition lock, a batchlog write and one or more remote writes per view. In 4.0 the concurrent_materialized_view_writes setting, 32 by default, bounds concurrent view writes, and the configuration comment notes that because a read is involved it should be limited by the lesser of the concurrent read and write limits. Write latency and throughput both move in the wrong direction, and they get worse with each additional view on the same base.

Changing a view key column writes a tombstone into the view each time, so a frequently changing key such as a status column builds tombstones in the view's partitions; see tombstones for why that hurts reads. Hot base partitions serialise on the partition lock. And because a view partition is keyed by the new column, a low-cardinality key such as country produces very large view partitions.

Detecting divergence

If you run views, check them. The sampler below reads random base token ranges and confirms each row is present in the view. It proves only that expected rows exist; a stale extra row needs the reverse check, scanning the view and confirming each row against the base.

# Base-to-view consistency sampler (DataStax Python driver). Run off-peak, throttled.
from cassandra.cluster import Cluster
from cassandra import ConsistencyLevel
from cassandra.query import SimpleStatement
import random

session = Cluster(["10.0.0.11"]).connect("app")
lookup = session.prepare("SELECT user_id FROM users_by_email WHERE email = ? AND user_id = ?")
lookup.consistency_level = ConsistencyLevel.ALL

def sample_range(start, end, limit=2000):
    q = SimpleStatement(
        "SELECT user_id, email FROM users WHERE token(user_id) > %s AND token(user_id) <= %s LIMIT %s",
        consistency_level=ConsistencyLevel.ALL, fetch_size=500)
    missing = []
    for row in session.execute(q, (start, end, limit)):
        if row.email is None:
            continue                                  # not in the view by definition
        if session.execute(lookup, (row.email, row.user_id)).one() is None:
            missing.append(row.user_id)
    return missing

# Pick random token ranges each night; alert if the missing rate is above zero,
# then rewrite affected base rows (same values, new timestamp) to regenerate view rows.

The usual repair for a divergent row is to rewrite the base row with its current values, which regenerates the view delta from a now-consistent base. For widespread divergence, dropping and rebuilding the view is often the only reliable fix, which is itself an operation measured in hours.

Alternatives

ApproachHow it stays in stepBest for
Application-maintained tableApplication writes both tables, often in a logged batch; see batches.Most production cases; you own the logic and can reconcile.
Storage-attached indexIndex built into each SSTable and memtable on the base replicas; see SAI.Filtering by a column within bounded partitions or with a partition key.
Change data capture projectionA consumer reads CDC and writes a derived table.Derived tables that can lag and need transformation.
Materialized viewServer-maintained, asynchronous, no base-view repair.Low-stakes lookups where you accept and monitor divergence.

The application-maintained table has the same asynchronous character if you write the two tables separately, but a logged batch guarantees that both writes eventually apply, and the logic is in code you can test, version and reconcile.

Operating views if you choose them

  • Enable deliberately: the 4.0 configuration uses enable_materialized_views and the 5.0 configuration uses materialized_views_enabled, both false by default; check the key for your version.
  • Cap the number of views: the 5.0 guardrails materialized_views_per_table_warn_threshold and materialized_views_per_table_fail_threshold (disabled by default at -1) stop a table from gathering views.
  • Repair the base on schedule and run a divergence sampler with alerts.
  • Choose view keys with high cardinality that rarely change.
  • Watch view write latency and batchlog replay activity alongside base write latency.

What to do next

  1. Inventory existing views: base table, view key, write rate and whether readers depend on exact results.
  2. For each view that feeds correctness-critical reads, plan a migration to an application-maintained table or SAI.
  3. For views you keep, run a nightly divergence sampler and alert on any missing rows.
  4. Confirm base-table repair runs within gc_grace_seconds on every table with views.
  5. Set the per-table view guardrails and leave view creation disabled on clusters that do not need it.
  6. Before creating any new view, test the view build time and write latency on production-sized data.
Key takeaway: A materialized view is a table that each base replica updates asynchronously after a lock, a local read and a batchlog write, sending to a paired view replica chosen by position. Client consistency levels cover the base only, repair never compares base with view, and a stale base replica produces a stale view row that nothing will notice. Use views only where you can tolerate and monitor divergence; for everything else maintain the second table in the application or use SAI, and keep base repair and a divergence check running if you keep views at all.