All articles

Apache Cassandra 6.0 Part 3 - Performance Optimisations

Performance Optimisations

I started using Cassandra in 2008 when sustained write throughput and horizontal scaling were a large part of its appeal. You could keep accepting writes while adding nodes or losing one and that made it a strong choice for workloads that needed to keep growing.

Other databases have since improved performance through different storage engines and approaches to scheduling with ScyllaDB becoming a prominent competitor. Comparing them means running the workload you care about on equivalent hardware with the same consistency requirements. The schema and data distribution can change the result as much as the balance of reads and writes or the operational limits you set.

Cassandra 5.0 and the work in 6.0 address several of the costs behind those comparisons through changes to allocation and busy code paths. You can also see improvements aimed at flushing and compaction where the work continues long after an application has received its write acknowledgement. This post follows those paths to explain where you might see lower resource use and what to measure on your own cluster.

This post describes pre-release Cassandra 6.0 behaviour as reviewed against 6.0-alpha3 rather than a finished production release. Check the release notes and upgrade documentation for the version you’re testing before using these settings.

Cassandra 5.0 and 6.0

Cassandra 5.0 added direct I/O for commit logs through Java native APIs and introduced a trie-based memtable implementation. It also brought the trie-indexed bti SSTable format and Unified Compaction Strategy so there were already new options for lookups and compaction. The 6.0 work continues into what happens when a frozen memtable is flushed or compaction reads an SSTable. It also reduces some of the work involved in serializing responses to read requests.

AreaCassandra 5.0 baselineCassandra 6.0 changes
Durable write I/ODirect I/O became available for commit log files.Direct I/O reaches compaction reads for compressed SSTables, where large sequential scans can otherwise displace query data from the page cache.
Memtables and flushingTrieMemtable improved memory use, GC efficiency, and lookup performance while data is in memory.The flush path is specialized to remove allocations, repeated checks, mapping work, and megamorphic calls as the frozen memtable becomes an SSTable.
SSTable lookups and compactionThe bti format and Unified Compaction Strategy gave 5.0 new storage and compaction options.Read, protocol, metadata, and response-serialization paths receive allocation and CPU reductions around those storage operations.
GC loggingGC work and its logging remain visible under sustained allocation.Async GC logging is enabled on JDKs that support it so log-file I/O is less likely to stall application threads.

To find out how these changes compare with ScyllaDB you’ll need to run both against the application’s workload. The changes below address identifiable CPU and allocation costs as well as page-cache pressure but they don’t establish a performance result for every cluster.

The following tickets give you the implementation details for the changes we’ll work through in this post.

Direct I/O for compaction reads

Start with a node serving reads from data that usually fits in the operating system’s page cache. Cassandra 5.0 already allowed direct I/O for the commit log but compaction still read its input SSTables through that cache. Pages retained for ordinary queries could then compete with a large sequential scan of SSTables being merged and replaced.

If the compaction reads more data than the available cache can hold it can evict pages that were serving your queries. Cassandra has to fetch that query data again while the device is already handling compaction and the kernel is reclaiming pages and writing dirty output. A large STCS compaction overlapping a frequently read dataset makes this behaviour easier to see although other compaction strategies can create the same pressure.

CASSANDRA-19987 lets compaction read compressed SSTables through direct I/O so those input bytes bypass the page cache. Uncompressed SSTables still use buffered reads and the merge continues to perform decompression and reconciliation before writing its output. Direct I/O also requires suitable alignment and buffering which is why the implementation depended on Cassandra’s internal buffering support.

You can select that path with the following setting while leaving ordinary client reads on their existing access path.

# cassandra.yaml
# auto inherits from disk_access_mode
# direct bypasses the OS page cache for compaction reads
compaction_read_disk_access_mode: auto

With compaction_read_disk_access_mode set to direct compaction can stop displacing cached query data with its sequential input scan. You’ll still need to watch disk contention because the same storage device has to serve both compaction and query I/O.

Jon Haddad reported CASSANDRA-19987 and Sam Lightfoot implemented the change with Ariel Weisberg and Maxwell Guo credited for review in the Apache Cassandra PR. I’m grateful for both the implementation and the work explaining how to evaluate it. Sam’s Direct I/O for Cassandra Compaction: Cutting p99 Read Latency by 5x walks through the page-cache problem and the results from his test setup.

Try buffered and direct compaction reads with the same schema and SSTable layout on the storage you intend to use. Keep application traffic and compaction concurrency unchanged while comparing p95 and p99 query latency with device queues and disk latency. Record cache hits and major faults together with compaction throughput and the time the query data takes to return to cache afterwards. Some devices may deliver better query latency at the cost of compaction throughput so check whether the node can still complete its maintenance work.

CASSANDRA-21134 extends direct I/O to compressed SSTable writes during background operations including compaction and streaming. It also covers cleanup and repair as well as upgrades while memtable flushes stay buffered so recently flushed data can benefit from the page cache. Part 7 works through the configuration and storage tests for that write path in more detail.

Memtable flush optimization

TrieMemtable in Cassandra 5.0 improved memory use and lookups before a flush while reducing garbage-collection work. Once the memtable is frozen Cassandra still has to walk its partitions and rows and serialize their cells into an SSTable. It updates metadata and builds indexes as it writes so a busy node can spend substantial CPU time flushing data.

Consider a write-heavy node where new mutations fill memtable memory as quickly as flushing can release it. Any delay in flushing can lead to backpressure or memory pressure before the application reaches the write throughput it needs. Flush threads also share CPU with mutation processing so you can encounter that limit before using all the device’s disk bandwidth.

Dmitry Konstantinov reported and implemented CASSANDRA-21083 to reduce the work performed inside that flush path. The table below follows the checks and allocations that can be avoided while writing the same frozen memtable.

OptimisationCPU and allocation savings
Update MetadataCollector clustering values only for the first and last clustering key in a partition.The SSTable metadata needs partition bounds, not the same update for every row written in the partition.
Split Cell.Serializer and MetadataCollector.update(Cell) call sites.A call site that sees many different concrete cell types becomes megamorphic, which prevents the JIT from making the inlining decisions available to a monomorphic or lower-polymorphism call site.
Precalculate counter-column information and move guardrail checks outside the per-row loop.The common flush path avoids repeated type checks, guardrail lookups, and hidden boxing for logging parameters.
Return early from row and deletion checks when no complex deletion or tombstone work is present.A normal live row does not pay for iteration or metadata checks required only by less common deletion cases.
Reduce NativeClustering serialization allocation and avoid remapping columns when the mapping has not changed.Serialization performs less object creation and avoids rebuilding data that is already in the expected form.
Use a flush iterator without column filtering.A flush writes the full frozen memtable, so it does not need the general filtering machinery used by other readers.

A busy flush can perform Cell operations millions of times per second so small costs repeat across a large amount of work. If a call site sees several unrelated concrete implementations the JIT has fewer opportunities to inline the calls. Splitting that call site gives it a more stable type profile and can reduce the cost each time a flush thread writes a cell.

Dmitry credits Branimir Lambov for reviewing the flush changes and the supporting profiles and benchmarks help explain the decisions in the patch. That review and regression testing are also what let us evaluate the optimization without losing sight of whether the output remains correct.

Run the comparison under sustained write load so the test reaches the point where flushing has to keep up with incoming mutations. Keep the memtable implementation and flush writer count unchanged along with compression and the data model on the same device. Compare flush duration and bytes written per second with write throughput and latency while watching allocator stalls and pending work. Heap allocation and flush-thread CPU help explain any change in GC activity and you’ll want to check that a faster write rate doesn’t leave compaction falling further behind.

Read-path optimisations

A read still involves work beyond finding data in a trie memtable or a bti SSTable. The coordinator plans the request and replicas read and merge rows before selection and serialization produce the response for native transport. Allocations along that path add young-generation collection work even when the amount of retained heap stays much the same.

The following changes reduce allocations and repeated calculations at specific points along the read and response paths.

TicketChangeExpected benefit
CASSANDRA-21199Allocation improvements in ProtocolVersion, StorageProxy, and MerkleTree.Lower per-request and metadata housekeeping allocation under coordinator and repair-related activity.
CASSANDRA-21360Removes allocations from miscellaneous read-path locations.Lower allocation rate under the query patterns that exercise those paths.
CASSANDRA-21362Avoids wrapping ByteBuffer values in cql3.selection.Selector.InputRow.Fewer short-lived wrapper objects while processing selected values.
CASSANDRA-21414Reduces the cost of calculating BTreeRow.minDeletionTime.Less CPU spent on row metadata that is examined repeatedly by read and storage paths.
CASSANDRA-21285Uses a lightweight moving average to size LocalDataResponse output buffers before row serialization.Fewer buffer resizes, allocations, and array copies for responses that exceed the old initial buffer size.

Take a response that’s larger than the output buffer initially allocated to hold it. Growing that buffer repeatedly means allocating a larger array and copying the existing content each time it runs out of space. Cassandra 6.0 uses a lightweight estimate of recent response sizes to start closer to the capacity it needs. The client receives the same CQL result while the server creates fewer intermediate buffers to produce it.

Use the application’s result sizes and paging settings when comparing 5.0 and 6.0 and include wide rows and selected collections. Keep the consistency levels unchanged and look at Java Flight Recorder allocation profiles for native transport and coordinator activity. Compare those allocations with request throughput and p95 and p99 latency while recording CPU use and GC pauses or concurrent-GC CPU. You’ll then be able to check whether fewer allocated bytes actually reduce the work competing with requests.

GC logging

CASSANDRA-21372 enables asynchronous GC logging on JDKs that support it so writing a GC event is less likely to stall on log-file I/O. You can check for fewer logging-related stalls while measuring the collector separately through allocation rate and pause time. Concurrent-GC CPU should also be recorded because changing how logs are written doesn’t change how much memory the collector needs to process.

Contributors

Jon Haddad reported the direct I/O issue and Sam Lightfoot implemented it with review from Ariel Weisberg and Maxwell Guo. Dmitry Konstantinov reported and implemented the flush work with Branimir Lambov credited for its review. C. Scott Andreas authored the LocalDataResponse allocation work with Caleb Rackliffe also credited as a co-author and both Caleb Rackliffe and Dmitry Konstantinov as reviewers.

These credits cover the contributions explicitly identified on the tickets and there is much more work behind the release. People maintain benchmarks and CI as well as the test infrastructure used to catch regressions before users encounter them. I’m grateful to everyone preparing releases and documenting behaviour alongside those helping users investigate problems on their own clusters.

Cassandra 5.0 and 6.0 test plan

Keep both versions on the same hardware and JDK with matching heap and replication settings. Use the same schema and compaction strategy with equivalent data and client traffic so the comparison exercises the changes described here. Check cache state too because a warm run on one version and a cold run on the other can obscure the effect you’re trying to measure.

For the before-and-after comparison I recommend using AxonOps to collect all Cassandra metrics every 5 seconds from every node. With table and keyspace metrics alongside JVM and operating-system data you can check whether a lower CPU figure came with more SSTables per read or a growing compaction backlog. Five-second samples make short bursts of latency or queued work easier to compare between runs when longer averages would smooth them out.

TestCassandra 5.0 baselineCassandra 6.0 comparisonMeasurements
Compaction with a live hot read setBuffered compaction reads, using the normal page cache path.Repeat with compaction_read_disk_access_mode: direct.p95/p99 read latency, major faults, page-cache activity, device queue depth, read/write latency, compaction throughput, and time to recover the hot set.
Sustained write and flushSame memtable implementation and memtable_flush_writers setting.Repeat after 6.0 flush-path changes.Write throughput and latency, allocator stalls, flush duration, flush CPU time, allocation rate, GC activity, pending flushes, and compaction backlog.
Read response serializationA result size and paging profile that represents the application.Repeat with 6.0 response-buffer and selection allocation changes.Native transport allocation rate, buffer-copy allocation, coordinator CPU, p95/p99 latency, response size, and throughput.
Mixed workloadReads, writes, repairs, and compaction at production-like concurrency.Repeat with the same traffic while enabling only the 6.0 setting under evaluation.Coordinator and replica latency separately, read amplification, SSTables per read, tombstones scanned, disk saturation, GC pauses, and concurrent-GC CPU.

Compare matching test durations at the same chart resolution and check your retention settings so the five-second data is still available when you review both runs. Keep Java Flight Recorder allocation profiles with the logs and exact configuration for each run so you can investigate any change in CPU use or GC activity.

Our Datadog Cassandra metric comparison shows the table-level and operational measurements missing from the default integration and how those gaps complicate investigations. Joaquin Casares also tested collector overhead under load in Monitoring Cassandra: The Cost of Collecting Metrics using Cassandra 4.1.3.

Series

This post is part of the Cassandra 6.0 series.

Sources

All articles