Transactional Cluster Metadata (TCM) and Cluster Metadata Service (CMS) in Cassandra 6.0
When a node joins Cassandra the rest of the cluster needs to agree on how ownership changes while data is moving. Transactional Cluster Metadata (TCM) is the architecture introduced by CEP-21 to order those changes along with other critical cluster metadata. Within it the Cluster Metadata Service (CMS) uses a small quorum of Cassandra nodes to accept and publish metadata changes in an agreed order.
Cassandra continues to use gossip for liveness and transient information while TCM records correctness-critical metadata in a linearized event log. Topology and token ownership belong in that log together with schema and the data placements used to decide which replicas can serve a request.
During bootstrap or decommission you need the old and new replica groups to remain consistent while ownership changes. Replacement and token movement need similar coordination and schema changes have their own dependencies between definitions. TCM assigns those changes an ordered history and an epoch with precomputed placements and checks that affected replica groups have observed one step before the next proceeds.
The 6.0 implementation also includes tools for inspecting and recovering metadata together with CMS rediscovery after addresses change. Peer-table repair and topology-safe changes to a live node’s data centre or rack are part of the same work.
This post follows pre-release Cassandra 6.0 as reviewed against 6.0-alpha3 and the details may change before GA. Check the release notes and upgrade documentation for the version you’re testing before planning a topology change.
Why Cassandra moved metadata out of gossip
I’ve been operating Cassandra for almost twenty years and schema disagreement has been one of the most painful problems to investigate and recover from. Nodes could catch up after missing a change but eventual propagation gave them no authoritative order for validating changes across the cluster.
Cassandra 4.x used timestamped last-write-wins schema changes that separate coordinators could accept against different local views. CEP-21 describes a table being created with a user-defined type while another coordinator accepts DROP TYPE and leaves that table unable to load after restart. Pulling or pushing schema could trigger migration storms without establishing an order for validating dependent changes. TCM validates each proposal against the latest committed metadata before accepting it into an ordered history.
With gossip-based topology updates a coordinator could still read from {A,B} while another wrote to {C,X} during range movement. CEP-21 shows how those non-overlapping replica sets could let a read miss a successful write and how some timing windows could leave writes only on replicas about to stop owning the range. TCM orders these transitions with separate read placements and write placements so temporary extra write replicas protect writes while ownership changes.
Cassandra’s masterless design has been valuable to me because normal reads and writes stay distributed across coordinators and replicas. Coordinating ownership changes between nodes with different information is difficult and I’m looking forward to TCM giving those changes a durable history without putting a central master into normal data requests.
I’m grateful to the committers and wider contributor community for the implementation and review as well as testing and compatibility checks. The documentation and release preparation take sustained effort too and I appreciate the people keeping users informed as the project works through a change of this size.
TCM architecture
Follow a metadata change from the node proposing it to a CMS member in the diagram below. CMS members are Cassandra nodes that validate the proposal against the latest metadata and reject conflicts with in-flight operations or replication invariants. A successful proposal becomes an immutable log entry committed through a CMS quorum and receives a monotonically increasing epoch.
CMS propagates the committed entry to the cluster and each node applies the ordered events to its local immutable ClusterMetadata view. Committing an event doesn’t require every node to receive it synchronously because propagation happens separately. The topology-operation rules control when the next step can proceed by requiring acknowledgements of the prior epoch from the affected replica groups.
The CMS members run inside Cassandra so you don’t need a separate ZooKeeper or etcd service to provide metadata consensus. Normal reads and writes still use Cassandra’s coordinator and replica path with a local placement snapshot rather than passing through CMS.
A coordinator includes the relevant epoch in internode messages while using its local placement snapshot for a read or write. If a replica is behind it can catch up before processing the request at that epoch. A replica that already knows a newer epoch sends that information back so the coordinator can catch up and validate its replica plan again.
Those epoch checks connect asynchronous metadata propagation to the replica plan used for an actual request. An older coordinator learns that its placement has been superseded before it can quietly complete a quorum request against an incompatible view.
Epochs, metadata log entries, and placement snapshots
You can read an epoch as a position in the committed metadata history identifying one immutable cluster state. It is an ordering value rather than a wall-clock timestamp so comparing epochs tells you which metadata a node has applied.
Cassandra precomputes read and write placements for a topology epoch so the coordinator can use a stable replica plan without recalculating placement on every request. When a request encounters a newer epoch the node can retrieve missing log entries and update its local snapshot. It then checks whether the original plan still satisfies the requested consistency level with that metadata.
The table below separates agreement on a metadata event from its later propagation to other nodes. CMS quorum consensus commits the event while each Cassandra node catches up with the resulting ordered history asynchronously.
| Stage | What is ordered or checked | Why it exists |
|---|---|---|
| Event submission | The CMS validates schema, topology, ownership, and in-flight operations against the current metadata. | A decommission that breaks a required replication factor, or a conflicting range movement, can be rejected before it changes placement. |
| CMS commit | A quorum of CMS members appends the event and assigns the next epoch. | There is one durable order for metadata changes. |
| Metadata publication | Nodes receive the new log entry and derive an immutable ClusterMetadata snapshot. | Nodes converge on the same directory, schema, and placements without relying on a race between gossip messages. |
| Topology-step gate | A majority of the relevant pre-change and post-change replica group acknowledges the previous epoch. | Read and write quorums retain an overlap while ownership is changing. |
| Request-time epoch exchange | Coordinators and replicas compare the epoch carried with a mutation or read request. | A stale replica plan can be caught up, rejected, or retried rather than being accepted silently. |
How TCM bootstraps a node without breaking quorum intersection
Let’s follow the four-step bootstrap example from CEP-21 using RF=2 with existing nodes A and B alongside C. Node X joins at token 150 and splits the range between tokens 100 and 200.
Initially the range (0,100] belongs to {A,B} and (100,200] belongs to {B,C} before the new node joins. Cassandra first splits the second range at 150 without changing ownership and then updates write placement before read placement. Superseded write replicas remain in place until the required nodes have acknowledged the preceding state.
Read the steps in order to see why the new node starts receiving writes before it can serve reads.
- Cassandra splits the ranges for the new token while keeping their existing read and write placements unchanged.
- It adds
Xto the write placements for the ranges it will own so writes reach both the old replicas andXwhile reads still use the old groups. Xstreams the required data before Cassandra can add it to the read placement and remove the outgoing replica from that read placement.- Cassandra removes the outgoing write replica only after the required acknowledgements show that a coordinator using the old read placement must learn about the newer epoch.
You can see a period when the outgoing replica still receives writes even though it has left the read placement. That overlap lets a coordinator using the old view and another using the new view retain the required quorum intersection. Removing the old write replica too early could let a read use replicas that no longer receive the new writes.
Decommission follows the reverse transition by adding recipient nodes to write placements before streaming from the leaving node. Cassandra then switches read placements and removes the departing write replica before merging equivalent adjacent ranges. Replacement and removal use corresponding metadata transitions along with token movement and other ownership operations.
Failure scenarios and recovery behaviour
If a topology change stops halfway through you’ll have a durable record of the steps that completed and the state still pending. The recovery procedure can use that record to resume or compensate for the change as appropriate for the operation.
| Scenario | What TCM records or detects | Operational consequence |
|---|---|---|
| A coordinator has an old epoch | A replica response carries a newer epoch than the coordinator used for its plan. | The coordinator catches up from the metadata log and checks the plan again. If the plan is no longer sufficient for the requested consistency level, the request fails rather than claiming success against an obsolete placement. |
A bootstrap stops after X becomes a write replica | The metadata log retains the pending bootstrap state and its epoch. | Cassandra does not let each node independently erase the pending state based only on local liveness. X must catch up and complete streaming before read placement can change; an unrecoverable operation needs the version-specific recovery procedure. |
| Two topology changes overlap | The CMS evaluates an event against in-flight range movements. | Disjoint operations may proceed concurrently. Operations that would interfere can be rejected until the first movement completes. |
| A node is down while metadata advances | The node’s persisted metadata is behind the log tail. | On restart it discovers the CMS or a peer, replays the missing events, rebuilds its local metadata snapshot, and then continues from the current state. |
| A node has a stale schema view | The coordinator or replica identifies a newer epoch during internode communication and can replay the ordered metadata entries. | Cassandra does not rely on independently timed schema pulls and pushes to decide the DDL order. A request that cannot satisfy its consistency requirements with the current metadata fails rather than reporting a success against a conflicting view. |
| CMS member addresses change while a node is down | A persisted address list can be stale even when CMS membership is still valid. | CASSANDRA-20476 adds a discovery protocol that builds a temporary address lookup, contacts the current CMS, and allows the agreed address change to be committed and disseminated. |
| A CMS quorum is unavailable | The CMS cannot commit a new metadata event. | Existing stable data traffic has its normal Cassandra behaviour, but new schema or topology changes that need CMS consensus must wait for CMS recovery. Treat this as a cluster-control-plane incident with a version-specific recovery runbook. |
A paused bootstrap can retain its pending state while CMS is still healthy enough to accept other valid metadata events. Losing CMS quorum prevents agreement on the next event and needs a different recovery procedure. Follow the procedure published for your Cassandra version and the actual metadata-log state because generic resets or ad hoc edits can make either situation harder to recover.
What Cassandra 6.0 adds around TCM operations
TCM first ships in Cassandra 6.0 after being developed on trunk when that unreleased line was numbered 5.1. The following changes are part of that same release line and cover operation through restarts and address changes as well as recovery and topology maintenance.
- CASSANDRA-20476 lets a restarting node rediscover CMS members after their addresses change. Node-ID based membership allows it to find current endpoints before catching up with metadata.
- CASSANDRA-20528 brings live-node data-centre and rack changes under the metadata-driven topology rules so location changes can be coordinated safely.
- CASSANDRA-19151 provides an offline cluster metadata tool for inspection and recovery when you need to examine state outside a running node.
- CASSANDRA-20525 adds a
nodetoolcommand for inspecting the cluster metadata log and directory virtual tables during an investigation. - CASSANDRA-21187 adds tooling to repair inconsistent information in
system.peersandsystem.peers_v2after checking the recorded state.
Use these tools with the documented procedure for the operation you’re investigating and retain the metadata state before changing it. They make the log and peer information easier to inspect alongside the node logs and gossip output available during recovery.
Initializing CMS after an upgrade
Once every node in the rolling upgrade is running Cassandra 6.0 you must run nodetool cms initialize on one node. That step is required to complete the TCM upgrade so include it in the upgrade procedure.
Schema changes are prohibited from the first restart into Cassandra 6.0 until CMS initialization finishes. Node bootstrap and decommission are also prohibited along with moves and replacements and assassinate during that period. Disable automation that could perform any of those operations before beginning the upgrade and keep it disabled until initialization completes.
CMS starts with a single member after initialization so expand it with nodetool cms reconfigure before treating the cluster as ready. Cassandra recommends at least three CMS members and three to seven per data centre distributed across failure domains. Nodes can catch up through metadata-log entries or snapshots with nodetool cms snapshot available for working with metadata snapshots.
Scaling TCM and the CMS
Only the CMS subset participates in consensus on metadata events so adding data nodes doesn’t require every node to join that consensus group. Ordinary reads and writes continue through their coordinators and replicas using local placement information.
A bootstrap still needs to stream data while existing nodes may receive extra writes during the placement transition. The next stage waits for the required epoch acknowledgements from affected replica groups so metadata propagation also contributes to progress. Schedule large-cluster operations with those limits in mind instead of assuming all bootstraps or decommissions can run together.
CEP-21 allows concurrent operations on disjoint replica groups when the replication invariants remain satisfied. To evaluate that concurrency you’ll need the affected ranges and available streaming bandwidth together with disk headroom and compaction load. Progress also depends on CMS and the participating nodes committing and acknowledging each epoch.
Choose CMS members across the failure domains the cluster needs to tolerate and check that maintenance leaves a quorum available. The placement decision governs the availability of metadata changes separately from the replica placement used for ordinary data. Adding every Cassandra node to CMS enlarges the consensus group without helping normal reads while too few members can leave topology or schema changes unavailable after a failure.
Benefits and constraints
| Benefit | What it changes | Constraint to keep in view |
|---|---|---|
| Ordered metadata history | Every committed topology or schema event has an epoch and immutable position in the log. | Operators need to understand whether a node is behind, whether an event is pending, and whether the CMS still has quorum. |
| Safer range movement | Read and write placements change in staged transitions with majority acknowledgement gates. | Bootstrap and decommission still need streaming capacity, disk headroom, and monitoring. |
| Better request-time safety | Epoch exchange exposes a coordinator or replica using an obsolete placement. | A request may fail during a metadata race rather than returning a false success. Clients should retain normal retry discipline. |
| Resumable operations | Pending metadata state is durable rather than being inferred independently from local liveness. | A permanently failed operation may require a deliberate cancellation or recovery procedure. |
| More inspectable control plane | Metadata log, directory state, and CMS tooling provide concrete evidence during diagnosis. | Tool output needs to be captured with the change record and compared to the intended placement. |
TCM operational checks
Before a topology change record the metadata epoch and CMS membership with the keyspace replication settings and intended placements. Keep the relevant system.peers state with that record and follow streaming and disk headroom while the operation runs. Compaction load and metadata-log events help explain delays alongside the epoch acknowledgements needed for each transition. Retain the final metadata view so you can compare the completed operation with the plan.
In AxonOps you can compare those records with node configuration and logs alongside disk activity and repair state. Compaction backlog and latency history help you check whether data movement affected requests while the topology operation was progressing. Seeing a node reach UN is useful but the surrounding history lets you examine how the cluster reached that state.
Series
- Part 1: Notes from Using Cassandra Since 2008
- Part 2: Accord Transactions
- Part 3: Performance Optimisations
- Part 4: Repair, Guardrails, and Observability
- Part 5: Zstd Dictionary Compression
- Part 7: Cursor Compaction and SSTable Writes
- Part 8: Storage-Attached Indexing and Schema Constraints
- Part 9: JDK 21 and Generational ZGC
- Part 10: Upgrade and Production Validation
Sources
- CEP-21: Transactional Cluster Metadata
- Apache Cassandra 6.0 CHANGES.txt
- CASSANDRA-20476: CMS member rediscovery and recovery
- CASSANDRA-20528: topology-safe data-centre and rack changes
- CASSANDRA-19151: offline cluster metadata tool
- CASSANDRA-20525: inspect cluster metadata log and directory
- CASSANDRA-21187: repair system peer tables
- Cassandra 6.0 ClusterMetadata source