Performance Tuning: Caching, Compaction Throttling, and JVM/GC Settings
Objective
Turn the storage-engine internals — memtables, SSTables, Bloom filters, and compaction — into a set of levers you actually pull when a cluster is slow. The companion concept on storage-engine internals explains what a memtable, an SSTable, a Bloom filter, and a compaction pass are, and why the LSM-tree design trades read complexity for write speed. This concept is about how to tune the three places where that trade becomes a dial you can turn: the read-side caches that sit in front of disk access, the resources given to the background compaction process that keeps SSTables from piling up, and the JVM heap and garbage collector that everything above ultimately runs inside of.
Use Cases
- A table is read-heavy and latency is dominated by disk seeks into SSTables — deciding whether a key cache (already on by default), a row cache, or leaning on the OS/chunk cache is the right lever, instead of reflexively reaching for more hardware.
nodetool compactionstatsshows pending compactions stacking up during a write-heavy period — deciding whether to raisecompaction_throughput_mb_per_sec, raiseconcurrent_compactors, or accept that the table is on the wrong compaction strategy (a data-modeling decision covered in the storage-engine internals concept, not a tuning knob).- A node is logging long GC pauses under
gc_warn_threshold_in_ms, and the question is whether to grow the heap, switch collectors, or both — and whether the book's Concurrent Mark Sweep (CMS) guidance is even runnable on the JDK version actually in production. - Comparing
nodetool proxyhistograms(coordinator-level read/write/range latency) againstnodetool tablehistogramson a specific table to decide whether a performance problem is cluster-wide (points at compaction, GC, or hardware) or table-specific (points at caching or compaction strategy for that one table). - Running
cassandra-stressor a customcqlstress-*.yamlworkload against a staging cluster before changing a productioncassandra.yamlorjvm.optionssetting, per the book's tuning methodology: change one parameter at a time and measure.
Deep Dive
Tuning methodology first: one change, one measurement
Before any specific knob, the book is explicit about the discipline this whole chapter assumes: "The suggested methodology for tuning Cassandra performance is to change one configuration parameter at a time and test the results. It is important to limit the amount of configuration changes you make when tuning so that you can clearly identify the impact of each change." Performance goals should be stated as both throughput and latency, at a percentile: "The cluster must support 30,000 read operations per second from the available_rooms_by_hotel_date table with a 99th percentile read latency of 5 ms" is the book's own example of a goal specific enough to tune against. Every setting below is a candidate change to make once, measure with nodetool tablehistograms/proxyhistograms or cassandra-stress, and either keep or revert.
Caching: key cache, row cache, chunk cache, counter cache
The storage-engine internals concept already covers what a Bloom filter is and calls it "a special kind of key cache" — a fast, in-memory, nondeterministic check that lets a read skip an SSTable's disk file entirely. The key cache is the literal cache that check feeds into. "Cassandra's key cache stores a map of partition keys to row index entries, facilitating faster read access into SSTables stored on disk." A Bloom filter says "maybe present, go check the file"; the key cache is what makes "go check the file" cheap once the Bloom filter says yes, by skipping straight to the row's index offset instead of scanning the SSTable's index from the start.
The book frames all caching decisions around three factors: "Consider your queries, and use the cache type that best fits your queries... Consider the ratio of your heap size to your cache size, and do not allow the cache to overwhelm your heap... Consider the size of your rows against the size of your keys. Typically keys will be much smaller than entire rows."
Key cache. Enabled by default per table — caching = {'keys': 'ALL', 'rows_per_partition': 'NONE'} is what DESCRIBE TABLE shows for a table that hasn't been touched. "Because the key cache greatly increases read performance without consuming a lot of additional memory, it is enabled by default." Its heap budget is shared cluster-wide via key_cache_size_in_mb, defaulting to "either 5% of the total JVM heap, or 100 MB, whichever is less." Disabling it (ALTER TABLE ... WITH caching = {'keys': 'NONE', ...}) is rarely the right move given how little memory it costs relative to the read-latency win.
Row cache. Caches entire rows rather than just an index pointer, at real memory cost, and the book is unusually blunt that it is often the wrong tool: "row caching tends to yield fewer benefits than key caching... in many cases, a row cache can yield impressive performance results for small data sets when all the rows are in memory, only to degrade on larger data sets when the data must be read from disk." The guidance is narrow: "Row caching is only recommended for read-heavy workloads, say 95% reads." Configure it per table via rows_per_partition (a count, or ALL; default NONE means disabled) — "If you are using a row cache for a given table, you will not need to use a key cache on it as well." The default implementation is off-heap (org.apache.cassandra.OHCProvider), which keeps a large row cache from directly pressuring GC.
Chunk cache. This is the one the book rates above row caching for most cases, precisely because it doesn't share the row cache's invalidation problem: "The chunk cache is considered to be more helpful than the row cache in most cases since it is more granular; the contents of the row cache are invalidated for an entire partition each time there is a write to the partition." The chunk cache stores decompressed chunks read off SSTable files (chunk size set per table by chunk_length_in_kb in the compression options), avoiding repeated decompression of hot, frequently-accessed data. It is enabled by default and sized via file_cache_size_in_mb.
Counter cache. Narrower purpose — "improves counter performance by reducing lock contention for the most frequently accessed counters" — with no per-table configuration, only a cluster-wide counter_cache_size_in_mb (default "2.5% of the total JVM heap, or 50 MB, whichever is less").
Operating the caches. All caches can be periodically persisted to disk under saved_caches and reloaded on startup to warm the cache without waiting for organic traffic — controlled per cache type by *_save_period and *_keys_to_save properties. At runtime, nodetool info reports hit rates per cache; nodetool invalidatekeycache / invalidaterowcache / invalidatecountercache clear a cache, and nodetool setcachecapacity / setcachekeystosave override the configured sizes without a restart — though "these settings will revert to the values set in the cassandra.yaml file on a node restart."
Compaction tuning: throughput, concurrency, and thresholds
The storage-engine internals concept covers why compaction exists and which strategy (STCS, LCS, TWCS, or 5.0's newer UnifiedCompactionStrategy) fits a given workload — that is a data-modeling decision made once per table. What follows here is the separate question of how much of the node's I/O and CPU budget compaction is allowed to consume once a strategy is chosen, since, as the book puts it, "compaction can be intensive in terms of I/O and CPU."
Throughput throttling. compaction_throughput_mb_per_sec in cassandra.yaml caps the rate compaction is allowed to write, checkable and settable live via nodetool getcompactionthroughput / setcompactionthroughput. "Setting this value to 0 disables throttling entirely, but the default value of 16 MBps is sufficient for most cases that are not write-intensive." This is a direct trade-off: throttle harder and client I/O keeps more headroom, but SSTables accumulate faster than they're merged; throttle less and compaction catches up faster at the cost of read/write latency for clients sharing the same disk.
Concurrency. If throughput alone doesn't clear a backlog, concurrent_compactors controls how many compaction threads run in parallel, configurable in cassandra.yaml or live via CompactionManagerMBean. It "defaults to the minimum of the number of disks and number of cores, with a minimum of 2 and a maximum of 8."
Per-table compaction thresholds. Distinct from the strategy itself, the threshold controls when a minor compaction triggers at all: "the number of SSTables that are in the queue to be compacted before a minor compaction is actually kicked off. By default, the minimum number is 4 and the maximum is 32." Set via CREATE TABLE/ALTER TABLE, or inspected and overridden per node with nodetool getcompactionthreshold / setcompactionthreshold. Too low and "Cassandra will spend time fighting with clients for resources to perform many frequent, unnecessary compactions"; too high and a single compaction event competes harder for resources when it finally runs.
Monitoring and control. nodetool compactionstats lists active compactions with a completed/total byte count and estimated remaining time — the first place to look when diagnosing a suspected compaction backlog. nodetool compactionhistory retains details of completed compactions. nodetool stop can halt a specific or all active compactions (Cassandra reschedules them), and nodetool disableautocompaction / enableautocompaction give per-keyspace/table control over whether compaction runs at all. A forced nodetool compact (major compaction) exists but the book flags real caveats: for SizeTieredCompactionStrategy use the -s flag "to request that Cassandra create multiple, smaller SSTable files rather than a single, large SSTable file," and if the actual goal is reclaiming space from deleted data, nodetool garbagecollect is the more targeted alternative to a full major compaction.
Testing a strategy change safely. Rather than changing a compaction strategy cluster-wide to see how it behaves, the book describes write survey mode: start a node with -Dcassandra.write_survey=true -Djoin_ring=false in jvm.options, configure the new CompactionParameters via JMX on that isolated node, then nodetool join it so it receives a live but non-authoritative stream of writes ("Writes to the test node place a minimal additional load on the cluster and do not count toward consistency levels"). Read performance for the new strategy is then tested by taking the node offline and querying it standalone.
JVM heap and GC tuning: from CMS toward G1GC
Heap sizing. Cassandra's default heap-sizing algorithm: on a machine with less than 1 GB RAM, heap is 50% of RAM; above 4 GB RAM, heap is 25% of RAM capped at 8 GB. Manual tuning uses -Xms/-Xmx, and the book's advice is to set them equal "to allow the entire heap to be locked in memory and not swapped out by the OS." -XX:+HeapDumpOnOutOfMemoryError is on by default in cassandra-env.sh and worth leaving on. As of the 3.0 release, GC- and heap-related flags moved into a dedicated jvm.options (now split further into jvm-server.options plus per-JDK-version files like jvm11-server.options), sourced by cassandra-env.sh.
The book's default: ParNew + CMS. The book describes the 3.x/4.0 default as two collectors working on different heap generations: the young generation uses the parallel copying collector (-XX:+UseParNewGC), tuned via -XX:SurvivorRatio (Cassandra sets 8, i.e. a 1:8 eden-to-survivor ratio, "fairly low, because the objects are living longer in the memtables") and -XX:MaxTenuringThreshold (Cassandra sets 1). The old generation uses Concurrent Mark Sweep, -XX:+UseConcMarkSweepGC — "this setting uses more RAM and CPU power to do frequent garbage collections while the application is running, in order to keep the GC pause time to a minimum." The book also notes a heap ceiling for this collector: "It is not recommended to set the heap larger than 12 GB if you are using the Concurrent Mark Sweep (CMS) garbage collector, as heap sizes larger than this value tend to lead to longer garbage collection pauses."
G1GC, per the book, as the pending alternative. "The Garbage-First garbage collector (also known as G1GC)... was intended to become the long-term replacement for the CMS garbage collection... G1GC generally requires fewer tuning decisions; the intended usage is that you need only define the min and max heap size and a pause time goal." The book is careful to note this wasn't a settled question as of the Revised 3rd Edition: "G1GC was originally the default for the Cassandra 3.0 release, but was backed out because it did not perform as well as the CMS for heap sizes smaller than 8 GB. The emerging consensus is that the G1GC performs well without tuning, but the default configuration of ParNew/CMS can result in shorter pauses when properly tuned." It also names two newer, more experimental collectors available on newer JDKs: ZGC ("its primary goals are to limit pause times to 10 ms or less, and to scale even to heaps in the multiple terabyte range") and Shenandoah ("works better than ZGC on smaller heap sizes... but SGC can have worse tail latencies"), both flagged as still maturing at the time of writing.
Book vs today
The collector recommendation has flipped, and CMS is no longer just "not preferred" — it is unusable on the JDK versions Cassandra now targets. The book's own text already hints this was in flux ("There has been considerable discussion in the Cassandra community about switching to G1GC as the default"), and it has since been settled:
- Checking the actual shipped configuration files confirms the change: Apache Cassandra 4.0's and 4.1's
conf/jvm11-server.optionsstill enable CMS by default (-XX:+UseConcMarkSweepGCactive, the G1 section commented out) — matching the book's description exactly. Apache Cassandra 5.0'sconf/jvm11-server.optionsandconf/jvm17-server.optionsflip this: G1GC is the active default (-XX:+UseG1GCuncommented, with tuned defaults like-XX:G1HeapRegionSize=16mand-XX:MaxGCPauseMillis=300), and the CMS block is present only as a commented-out legacy option. - This tracks a hard JDK fact, independent of Cassandra's own preference: CMS was deprecated in JDK 9 (JEP 291) and removed entirely in JDK 14 (JEP 363). Cassandra 5.0 supports running on JDK 11 or JDK 17 — meaning on a JDK 17 deployment, the book's CMS flags are not just non-default, they will fail to start the JVM at all. There is no
jvm17CMS configuration shipped because there is no CMS to configure on that JVM. - One inconsistency worth flagging rather than papering over: the prose "Hardware Choices" page in the current Apache Cassandra documentation still repeats older guidance — "heaps smaller than 12GB should consider ParNew/ConcurrentMarkSweep garbage collection" — even in the docs tree published alongside 5.0. The actual shipped
jvm11-server.options/jvm17-server.optionsdefaults have already moved past this advice to G1GC regardless of heap size. Treat the shippedconf/jvm*-server.optionsfile for your running version as the source of truth over prose documentation that may lag a release or two behind. - Practical takeaway: a team running Cassandra 5.0 out of the box is already on G1GC without having tuned anything, which is exactly the "fewer tuning decisions" property the book credits G1GC with. A team upgrading an older CMS-tuned cluster to a JDK where CMS no longer exists needs to migrate its
jvm.optionsdeliberately, not assume the old file still applies.
Trade-offs
- Caching trades heap (or off-heap) memory for disk I/O, and the three cache types are not interchangeable defaults. The key cache is nearly free and should almost always stay on; the row cache is a targeted tool for a narrow case (95%+ reads, hot small rows) and a bad default because "the wrong setting can easily lead to more performance issues than it solves"; the chunk cache benefits a wider range of read patterns because partition writes don't blow away its entire cached contents the way they do the row cache's. Sizing any of them too generously starves the JVM heap the rest of Cassandra depends on — the book's own warning to "not allow the cache to overwhelm your heap" is the whole caching story in one line.
- Compaction throttling is a direct latency-vs-catch-up trade, not a "set once" value. A low
compaction_throughput_mb_per_secprotects client-facing read/write latency during normal operation but risks a growing SSTable backlog under sustained write pressure — exactly the backlog that (per the storage-engine internals concept) degrades reads by forcing them to check more files. Raising throughput orconcurrent_compactorsclears the backlog faster at the direct cost of I/O and CPU that clients are also competing for. There is no throttle value that is correct at both light and heavy load; this is a setting to revisit under stress testing, not set once at cluster creation. - G1GC's "requires fewer tuning decisions" is a genuine advantage, not a guarantee it always wins. The book's own history — G1GC was Cassandra 3.0's default, then backed out for underperforming CMS below 8 GB heaps — shows collector performance depends on heap size and workload, not just release date. That Cassandra 5.0 now ships G1GC as default is strong evidence the collector and JVM have matured since that 3.0 backout, but it is not evidence that G1GC beats every possible tuned CMS configuration; it is evidence that G1GC beats untuned CMS for most workloads, which is what a shipped default should optimize for.
- The heap-size-equals-Xms-equals-Xmx recommendation trades startup flexibility for GC predictability. Locking min and max heap to the same value avoids the cost of the JVM growing the heap under load, but it also means a node cannot temporarily use less memory during quiet periods — the full heap commitment is permanent from process start, which matters when co-locating Cassandra with other memory-hungry processes on the same host.
- Off-heap memtable and row-cache storage reduces GC pressure at the cost of a more complex memory budget.
memtable_allocation_type'soffheap_buffersandoffheap_objects, and the row cache's default off-heapOHCProvider, all move data out of the space the garbage collector has to scan — directly helping whichever collector is in use, CMS or G1GC alike. The trade is that off-heap memory doesn't show up in heap-focused monitoring and still has to be budgeted against total system RAM alongside the heap itself; a node sized only by heap usage can still run out of memory.
Documentation Links
- Jeff Carpenter and Eben Hewitt, "Cassandra: The Definitive Guide", Revised 3rd Edition (O'Reilly, 2022) — Chapter 13, "Performance Tuning"
- Apache Cassandra Documentation — Hardware Choices (heap sizing and GC guidance)
- Apache Cassandra source — conf/jvm11-server.options (5.0 branch, G1GC enabled by default)
- Apache Cassandra source — conf/jvm17-server.options (5.0 branch)
- OpenJDK JEP 363 — Remove the Concurrent Mark Sweep (CMS) Garbage Collector
- Apache Cassandra Documentation — nodetool compactionstats