All articles

Apache Cassandra 6.0 Part 2 - Accord Transactions

Accord in Cassandra 6.0

Accord gives Cassandra 6.0 a way to coordinate transactions across multiple partitions through General Purpose Transactions under CEP-15. The related support for BEGIN TRANSACTION mutations that touch multiple partitions brings that coordination into CQL so we can work through an example using ordinary tables.

This post describes pre-release Cassandra 6.0 behaviour as it stood in 6.0-alpha3 when the article was updated. Check the release notes and upgrade documentation for the version you’re testing before relying on these details.

Accord is beta-designated and ships disabled through accord.enabled: false so you’ll need to enable it before trying the examples. Test your transaction workload and driver behaviour together with failure handling and the migration procedure before considering it for production.

I am grateful to the Cassandra committers and contributors working through the difficult parts of adding transactions to the database. Getting the CQL to run is only part of the work because the same transaction has to behave correctly when a node restarts or topology changes while requests are delayed.

Having used Cassandra since 2008 I want to understand how this coordination works and what it costs for the application. The data model still determines which partitions and replicas take part so it belongs in the discussion from the beginning.

Consider transferring a balance between two accounts stored in different partitions as in Phil Eaton’s Cassandra 6 exploration. You need the debit and credit to succeed together so two independent conditional writes would leave the application responsible for resolving a partial transfer. Paxos-based lightweight transactions provided conditional updates within Cassandra’s existing limits and Accord extends the options to general multi-partition transactions.

How Accord runs a transaction

CEP-15 describes a leaderless timestamp protocol that can coordinate transactions across any set of keys with strict-serializable isolation. You can follow the agreement steps without looking for a permanent cluster-wide leader because Accord doesn’t introduce one. Cassandra integrates the protocol with CQL and schema handling as well as the messaging and storage paths.

Before a transaction runs it declares the reads and conditions it needs together with the writes it may perform. Cassandra uses those declarations to find the participating partitions and ranges and assigns the transaction ID t0 as its initial timestamp. Knowing the data involved lets Accord check for conflicts before it executes the transaction.

StageWhat happens
PreAcceptThe coordinator sends the proposed transaction and its initial timestamp to replicas for every participating shard. Each replica records that it has seen the transaction, returns that timestamp when it can accept it, and includes lower-timestamp conflicting transactions as dependencies.
Fast pathWhen a fast-path quorum in every participating shard returns the initial timestamp, the coordinator commits the transaction with the combined dependency set. This is the normal one-round-trip agreement path described by CEP-15.
Slow pathA replica that has seen a newer conflicting timestamp responds with a higher timestamp. If the coordinator cannot take the fast path, it selects the highest timestamp it received and sends Accept to a simple quorum in each shard so that choice is durable before Commit.
ExecutionAfter Commit, the coordinator reads from the participating shards with the relevant dependencies attached. Replicas wait for earlier dependencies to commit and apply, return their read results, and then receive Apply for the final transaction result.
RecoveryA replica that has witnessed an incomplete transaction can coordinate recovery. It asks replicas for their local state, work it must wait for, and superseding transactions, gathers a simple quorum from every shard, then resumes the highest durable stage or selects the appropriate slow-path timestamp.

The coordinator is a Cassandra node and only the replicas responsible for the declared data participate in the transaction. Each participant records local transaction state so a replica that witnessed the work can coordinate recovery if the original coordinator stops making progress.

Fast path

Mermaid sequence diagram for Accord's fast path, showing Node 2, Node 4, and Node 6 receiving PreAccept, confirming the initial timestamp, then reading and applying the transaction result.

In the first diagram you can follow the initial timestamp through agreement before the transaction reads and applies its result. CEP-15 uses fast-path electorates so a minority of unavailable replicas doesn’t automatically force a slower agreement path or require every replica to reply.

Slow path

Mermaid sequence diagram for Accord's slow path, showing Node 6 reporting a higher timestamp, a simple quorum accepting the selected timestamp, then Commit, reads, and Apply.

If the responses don’t support the initial timestamp the coordinator chooses the highest timestamp it received and records that choice through Accept with a simple quorum. Once that step succeeds it can send Commit and carry out the execution reads before distributing the result through Apply.

Recovery after a replica failure

Mermaid sequence diagram showing Accord recovery after the original coordinator stops before a final decision and Node 4 is unavailable. Node 2, which witnessed the transaction, gathers state from Node 6 and resumes the strongest durable stage or proceeds through the slow path.

A replica can become unavailable while the transaction still completes through a fast electorate or the slow path as long as the required quorum remains available. The recovery diagram shows a different situation where the original coordinator stops before making a decision and leaves an incomplete transaction behind.

A participant that witnessed the transaction asks the other replicas what they recorded and which earlier work still needs to finish. Their responses identify progress through PreAccept and Accept as well as any later Commit or Apply state. Recovery continues from the strongest recorded stage when it finds a transaction that has already been accepted or progressed further.

If no stage has been decided the recovering participant checks whether the reports still support the initial timestamp or require a higher timestamp through the slow path. It waits and retries where earlier transactions need to finish before this transaction can continue.

Preparing tables for Accord

After enabling Accord in cassandra.yaml you’ll need to select a transactional mode for each participating table. The example below creates a new table with transactional_mode = 'full' so its operations can run fully through Accord.

CREATE TABLE commerce.orders (
  order_id uuid PRIMARY KEY,
  status text,
  sku text
) WITH transactional_mode = 'full';

CREATE TABLE commerce.allocations (
  sku text,
  order_id uuid,
  quantity int,
  PRIMARY KEY (sku, order_id)
) WITH transactional_mode = 'full';

You can also use mixed_reads to send writes and serial operations through Accord while keeping ordinary non-serial reads on Cassandra’s existing eventually consistent path. Moving an existing table requires a range-migration procedure that includes full repair and Paxos repair before the migration is complete. During phase two the first access to each key also incurs a Paxos repair round trip so allow for that work when testing the migration.

Counter tables are unsupported and you’ll also need to check the consistency levels your drivers use before enabling Accord. Reads accept ONE and QUORUM as well as SERIAL and ALL with ANY additionally accepted for writes. Requests using LOCAL_QUORUM are rejected along with those at TWO or THREE so test the levels your application actually sends. Once a table has fully migrated to transactional_mode = 'full' Cassandra ignores the supplied consistency level and commits through Accord with ANY semantics.

A cross-partition transaction

Let’s take an order that’s still pending and allocate it only if an allocation row doesn’t already exist. The transaction below checks both conditions before changing the order state and creating the allocation together. LET names the reads used in those checks and IF ... THEN controls the mutations for the whole transaction.

BEGIN TRANSACTION
  LET requested_order = (
    SELECT status FROM commerce.orders
    WHERE order_id = 7a3e2a9e-4d51-4c72-a0ee-000000000001
  );
  LET existing_allocation = (
    SELECT quantity FROM commerce.allocations
    WHERE sku = 'widget-42'
      AND order_id = 7a3e2a9e-4d51-4c72-a0ee-000000000001
  );
  SELECT requested_order.status;
  IF requested_order.status = 'PENDING' AND existing_allocation IS NULL THEN
    UPDATE commerce.orders
      SET status = 'ALLOCATED'
      WHERE order_id = 7a3e2a9e-4d51-4c72-a0ee-000000000001;
    INSERT INTO commerce.allocations (sku, order_id, quantity)
      VALUES ('widget-42', 7a3e2a9e-4d51-4c72-a0ee-000000000001, 3);
  END IF
COMMIT TRANSACTION;

Running the example lets you see the transaction boundary and the CQL used to coordinate those two changes. A complete inventory-reservation system would also need to check available stock and handle repeated requests without allocating it twice. Expiry and cancellation need their own behaviour together with a plan for handling failed or retried requests.

Keep the condition on the transaction block because Accord coordinates both the decision and the timestamp for the complete operation. Cassandra’s transaction tests reject per-statement IF conditions and custom USING TIMESTAMP clauses inside that block.

Strict-serializable isolation covers each transaction’s declared work so check the scope of the reads your application makes. Cassandra uses one transaction per page or sub-range for paged and partition-range reads which means you’ll need separate correctness tests for those operations.

Before Accord you would usually handle an operation spanning multiple partitions by choosing one of the following approaches.

  • Remodel the data so the operation fits a single partition
  • Coordinate the operation in the application
  • Use another system for that workflow

You can still choose any of those approaches when it fits the application and its performance requirements. Accord gives you another option when the operation needs cross-partition coordination and you’ve tested the cost of providing it.

Choosing transactional workloads

Start with an operation that already needs related records to change together and check how Accord could handle it. Moving a normal Cassandra write path to transactions adds coordination that the application may have no reason to pay for.

You’ll still need to understand the participating partitions and replicas because they determine how the transaction behaves when messages are delayed or a node fails. Test the consistency requirements and latency under those conditions before deciding whether the operation should use Accord.

The following cases let you evaluate coordination inside Cassandra before deciding whether to keep that work in an application service.

  • Update related records that must change together
  • Apply conditional state transitions across multiple partitions
  • Reduce application-side reconciliation for a specific workflow
  • Keep the coordination for a consistency-sensitive operation close to the data

Run those operations with the contention and retry behaviour you expect from the application before judging how well they fit.

Production considerations

The changelog records tail latency improvements alongside fixes to batch atomicity and transaction timestamp handling. It also covers clean shutdown and restart behaviour together with topology serializer improvements that affect how the protocol operates through cluster changes.

Include those situations in your tests because a successful transaction on an idle cluster tells you little about recovery during a restart. You’ll want to see how it behaves while the cluster changes and the application continues to submit work.

Give Accord-backed transactions their own workload test so you can examine the following measurements separately from ordinary Cassandra requests.

  • p95 and p99 transaction latency
  • Timeout rates
  • Invalid request rates
  • Retry behaviour
  • Contention patterns
  • Coordinator and replica latency separately
  • Behaviour during restart and node replacement as well as topology changes

Keep transactional writes visible in their own charts so you can compare a problem affecting them with the ordinary write workload.

Where AxonOps fits

For Accord you’ll need monitoring and alerts that let you follow the transaction workload alongside the rest of the cluster. A runbook should explain which measurements to compare when transactions slow down or stop completing.

During Cassandra 6.0 testing compare transaction latency with coordinator and replica latency to see whether the same delay affects ordinary requests. Check timeouts and unavailable exceptions alongside heap and disk pressure while keeping any workload-specific contention visible in the investigation.

Series

This post is part of the Cassandra 6.0 series.

Sources

All articles