A Cassandra cluster that ingests 20 MB/s of client data can easily write several hundred MB/s to its disks. None of that is a bug. Every byte is replicated, logged, flushed and then rewritten by compaction, sometimes many times, and occasionally streamed again by repair. This ratio of physical bytes written to logical bytes accepted is write amplification. It decides how many disks you need, how fast your SSDs wear out, and how much I/O is left for reads.

The general LSM theory, covering level sizing, the read/write/space trade-off and SSD endurance, is in Write amplification architecture. This article is the Cassandra version: where each multiplier comes from in a real cluster, how each compaction strategy changes it, how to measure it with the metrics Cassandra exposes, and which levers actually reduce it. Multipliers below are approximations for reasoning about capacity, not measured constants; measure your own cluster before acting on them.

Where the extra bytes come from

Follow one mutation through the write path, which is described in detail in Cassandra write and read path.

  1. Replication. The coordinator sends the mutation to every replica in every datacenter, so with replication factor 3 the cluster does three times the work below. This is usually the largest single multiplier, and it is the one you choose deliberately for durability.
  2. Commitlog. Each replica appends the mutation to its commitlog before acknowledging it, roughly one more write of the serialized mutation. Commitlog segments are recycled once their memtables are flushed, so the space is reclaimed but the bytes were still written.
  3. Flush. When a memtable fills or is flushed for another reason, it is written out once as an immutable SSTable. With table compression enabled, the default, the SSTable is usually smaller than the raw data. See memtable flush architecture.
  4. Compaction. SSTables are merged and rewritten to bound read cost and purge overwritten and deleted data. How many times each byte is rewritten depends almost entirely on the compaction strategy and the data volume per table per node.
  5. Repair, hints and streaming. Hints replay writes missed while a replica was down. Repair streams ranges whose Merkle trees disagree, and full repair can over-stream because a mismatch is detected per range, not per row. Bootstrap, decommission and rebuild stream whole token ranges. Every streamed SSTable then goes through compaction on the receiving node.
Where one client byte is written again in a Cassandra clusterClient write1 logical byteCoordinatorx RF replicasCommitlogappend, ~1xMemtableRAM, 0x diskFlushSSTable, ~1xCompactionSTCS/T4: ~levels; LCS/L10: much moreHintsreplayed after outagesRepair and streamingmismatched ranges, bootstrapDisk bytes written per nodecommitlog + flush + compaction + hints + streamingCluster write amplification = total disk bytes written on all replicas / logical bytes written by clientsReplication multiplies everything; compaction strategy decides the largest per-node term.
Replication multiplies every term; on each replica the commitlog and flush cost about one write each, and compaction usually dominates.

A useful per-node approximation, for data that lands on that node, is WA_node = commitlog + flush + compaction, and for the cluster WA_cluster = RF x WA_node + repair and hint overhead. The SSD's own internal amplification, from garbage collection inside the flash translation layer, sits underneath all of this and is invisible to Cassandra.

Indexes and views add their own terms. A materialized view turns each base-table write into extra mutations on the view's replicas, and the view table is then flushed and compacted like any other table, so it carries the full per-node cost a second time. Storage-attached indexes in Cassandra 5.0 are built per SSTable, so index components are written at flush and written again whenever compaction rewrites that SSTable. When you add a view or an index, count it as another table in the write budget, not as a free read optimisation.

How each compaction strategy changes the bill

Compaction is the term you control most. The strategies are compared in depth in Cassandra compaction strategies; here they are seen only through their write cost.

SizeTieredCompactionStrategy (STCS) merges SSTables of similar size once there are enough of them (the default min_threshold is 4). Each merge roughly quadruples SSTable size, so a byte is rewritten about once per tier, and the number of tiers grows with the logarithm of data size divided by flush size. With 100 MB flushes and 100 GB of table data on a node, that is about log4(1000), around 5 rewrites. STCS has the lowest write cost of the general-purpose strategies, paid for with more SSTables per read and temporary space for large merges.

LeveledCompactionStrategy (LCS) keeps fixed-size SSTables (160 MB by default) in levels that each hold about 10 times more than the previous one (fanout_size defaults to 10). Moving data from one level to the next merges it with the overlapping SSTables of the larger level, so each promotion can rewrite up to around fanout times the incoming data. Over several levels, rewrites per byte can reach tens. In exchange, most reads touch one SSTable per level, and space overhead stays low. When writes outpace it, level 0 backs up, and LCS falls back to size-tiered merging inside L0, which adds still more rewriting.

TimeWindowCompactionStrategy (TWCS) groups SSTables into time windows. It runs size-tiered compaction within the current window, compacts each window into one SSTable once it closes, and then leaves old windows alone. For time-series data written in time order and expired by TTL, whole SSTables are dropped instead of being rewritten, and write amplification is close to the minimum. Out-of-order writes, read repair and explicit deletes mix old data into new windows and undo much of that benefit.

UnifiedCompactionStrategy (UCS), added in Cassandra 5.0, expresses tiered and leveled behaviour with one parameter per level. scaling_parameters defaults to T4, which behaves like STCS with threshold 4, while values like L10 behave like leveled compaction with fanout 10. The documentation states write amplification as proportional to the number of levels for tiered settings and to (f - 1) times the number of levels for leveled settings. It also splits output into shards by token range, with a 1 GiB default target SSTable size, which keeps individual compactions smaller and more parallel.

StrategyRewrites per byte (rough)Best fitWrite cost risk
STCS / UCS T4about one per tier, often 4-8write-heavy, few overwriteslarge merges need temporary space
LCS / UCS L10can reach tensread-heavy, frequent updatesL0 backlog under heavy writes
TWCSabout 1-2 when data arrives in orderTTL time seriesout-of-order data and deletes

Measuring write amplification on a real cluster

Do not estimate when you can measure. Cassandra exposes cumulative counters per table on JMX under org.apache.cassandra.metrics:type=Table,keyspace=<ks>,scope=<table>: BytesFlushed (total bytes flushed since restart) and CompactionBytesWritten (total bytes written by compaction since restart). The compaction subsystem also exposes BytesCompacted. nodetool compactionhistory lists recent compactions with their input and output bytes, which shows which tables are being rewritten most.

Sample the counters twice, take the deltas, and compare them with the logical bytes your application wrote in the same interval, which you should count in the client or ingestion layer. Most metrics exporters can expose these JMX counters to your monitoring system.

def write_amplification(t0, t1, logical_bytes, rf):
    """t0/t1: dicts of per-table counters summed across all nodes at two times.
    logical_bytes: bytes accepted from clients in the same interval (before replication)."""
    flushed   = t1["BytesFlushed"] - t0["BytesFlushed"]
    compacted = t1["CompactionBytesWritten"] - t0["CompactionBytesWritten"]
    return {
        "replicated_logical": rf * logical_bytes,
        "compaction_per_flushed_byte": compacted / flushed if flushed else None,
        "sstable_wa_cluster": (flushed + compacted) / logical_bytes,
        # Add commitlog bytes and streaming from OS-level disk counters (for example
        # iostat on the commitlog and data devices) for the full physical figure.
    }

Two ratios matter most. compaction_per_flushed_byte tells you how many times compaction rewrites each flushed byte, which is the strategy's real cost on your data. The cluster ratio against logical bytes, cross-checked with device-level write counters, tells you whether disks and SSD endurance can sustain the load. Compare both over a full day, because compaction is bursty.

Worked example: sizing disks for endurance

A six-node cluster ingests 20 MB/s of logical writes into one table with RF 3, so each node accepts 20 x 3 / 6 = 10 MB/s. Assume compression halves the size on disk and use STCS with about 5 rewrites per byte, both of which you would verify with the counters above.

Term per nodeMB/sNote
Commitlog10uncompressed by default
Flush5compressed SSTables
Compaction, STCS255 rewrites of 5 MB/s
Total40about 3.5 TB/day

At about 3.5 TB written per node per day, a 3.84 TB drive rated for one drive write per day is already close to its endurance budget, before repair or a bootstrap. Switching the table to LCS because reads got slower might triple the compaction term. With an assumed 15 rewrites, compaction alone becomes 75 MB/s and the node writes 90 MB/s, about 7.8 TB/day, roughly twice the drive's rating. The same switch on a table that is mostly updates to a small working set might be the right call, because LCS discards overwritten data sooner. That is exactly why the decision needs the measured ratio rather than a rule of thumb.

Levers that actually reduce it

  • Match the strategy to the data. TTL time series in time order goes to TWCS. Append-mostly data that is rarely updated goes to STCS or UCS with tiered settings. Use LCS or leveled UCS settings only where read latency on updated rows justifies the extra writes.
  • Flush less often, in bigger pieces. Larger memtables produce larger initial SSTables, and with tiered strategies that removes the lowest tiers. Check heap or off-heap memtable space and the number of tables sharing it, because many small tables flush small SSTables.
  • Write less. Avoid rewriting unchanged columns, read-before-write patterns that rewrite whole rows, and deletes that could have been TTLs. Every tombstone is written, compacted and eventually purged; see Cassandra tombstones.
  • Use incremental repair where it fits. Repairing only unrepaired data, and keeping it separate from repaired SSTables, avoids re-streaming ranges that already agree.
  • Keep compression on. It reduces flush and compaction bytes for the CPU cost of compressing.
  • Do not mistake throttling for reduction. compaction_throughput spreads the same bytes over more time. Set too low, it builds a compaction backlog that hurts reads; it never lowers total bytes written.

Failure modes

  • Compaction backlog. Pending compactions climb steadily during peak hours because compaction cannot keep up with flushes. SSTable counts and read latency rise with them. Either the strategy amplifies too much for the disk bandwidth, or compaction is throttled too hard.
  • SSD wear-out ahead of plan. The drive's wear indicator falls faster than its warranty assumes, because capacity planning counted RF and flushes but not compaction.
  • TWCS windows that never settle. Late-arriving data or read repair writes old timestamps into new SSTables, and closed windows are recompacted again and again.
  • A major compaction on STCS. Forcing everything into one huge SSTable rewrites the whole table, and that SSTable then never finds similar-sized partners to merge with.
  • Repair storms. Full repair across many ranges at once streams a large share of the data, and the receiving nodes recompact all of it.

What to do next

  1. Export BytesFlushed and CompactionBytesWritten for your largest tables and compute compaction bytes per flushed byte over a full day.
  2. Count logical bytes written by your clients for the same day, and compute cluster write amplification including RF.
  3. Compare the daily bytes written per node with your SSDs' endurance rating, leaving headroom for repair and node replacement.
  4. For each large table, check that its compaction strategy matches its pattern: time-ordered TTL data, append-mostly data or update-heavy data.
  5. Graph pending compactions and SSTables per read next to the amplification ratio, so strategy changes show their cost in both directions.
  6. Before changing a strategy in production, test it on one node or in a staging copy with replayed traffic, and compare the measured ratios.
Key takeaway: Write amplification in Cassandra is replication times the per-node cost of commitlog, flush and compaction, plus repair and hints, and compaction strategy is the term you control most. Measure it from BytesFlushed and CompactionBytesWritten against the logical bytes your clients write, check the result against disk bandwidth and SSD endurance, and change strategy only when the measured ratio says the trade is worth it.