Hello! Welcome back.
In our last lesson, we explored Cassandra's storage engine, understanding how its Log-Structured Merge-Tree (LSM-Tree) architecture achieves phenomenal write throughput by creating immutable SSTables. We concluded by noting that this process necessitates a background task called compaction to manage the proliferation of these files.
Today, we will dive deep into that process. This lesson directly addresses the learning outcome: Configure compaction strategies in Cassandra for different workload patterns. Compaction is not a one-size-fits-all operation; it's a set of tunable strategies that critically impact your cluster's performance.
We will cover:
- The core trade-offs of compaction: Read, Write, and Space Amplification.
- The "legacy" compaction strategies: Size-Tiered (STCS), Leveled (LCS), and Time-Windowed (TWCS).
- The modern, recommended approach: Unified Compaction Strategy (UCS), and how to configure it for different workloads.
- Practical aspects of configuring and monitoring compaction.
Your experience with optimizing high-load systems has undoubtedly involved balancing competing resources—CPU, I/O, memory. Compaction is a prime example of this balancing act within a database, and understanding its mechanics is key to operating Cassandra effectively at scale.
1. The Core Problem: Amplification and Tombstones
Compaction's primary goal is to merge SSTables to improve read performance and reclaim disk space. However, this process itself consumes I/O and CPU. The choice of strategy is about managing the following trade-offs:
- Read Amplification (RA): The number of SSTables that must be accessed to satisfy a single read query. High RA leads to slower reads.
- Write Amplification (WA): The number of times a single piece of data is rewritten to disk due to compactions. High WA consumes more disk I/O and can reduce the lifespan of SSDs.
- Space Amplification (SA): The ratio of the total disk space used by SSTables to the actual size of the live data. High SA means you are storing a lot of stale or deleted data, wasting disk space.
A crucial aspect of space reclamation is handling deletes. As we've touched upon, Cassandra uses tombstones (deletion markers) instead of immediately removing data. These tombstones must be replicated to ensure deletes propagate through the cluster. Compaction is the process that eventually purges both the tombstone and the old data it marks for deletion.
Compaction | Apache Cassandra Documentation
To understand how deletes are managed, let's review the role of tombstones and the gc_grace_seconds setting. This is fundamental to configuring compaction, especially for workloads with frequent deletes or updates.
Please read the section titled "Tombstones and Garbage Collection (GC) Grace". Focus on why tombstones are necessary in a distributed system and the role gc_grace_seconds plays in preventing data resurrection.
As the documentation explains, a tombstone can only be permanently removed during compaction after it has existed for longer than gc_grace_seconds. This period (defaulting to 10 days) acts as a safety window, ensuring that a temporarily downed node has enough time to receive the delete marker via repair before it's purged from the cluster. If a node is down longer than this period, deleted data can be "resurrected" during a repair. An effective compaction strategy must be able to efficiently process and purge these tombstones.
2. The Legacy Strategies: STCS, LCS, and TWCS
Before the introduction of the Unified Compaction Strategy, operators had to choose one of three specialized strategies. Understanding them is valuable as they are still in use and their concepts form the basis of the modern approach.
Compaction | Apache Cassandra Documentation
The official Cassandra documentation provides a concise summary of the three main legacy strategies and their intended use cases.
Read the section titled "Strategies". It provides a brief description of Size Tiered, Leveled, and Time Window Compaction Strategies.
Let's break down the trade-offs for each:
| Strategy | Mechanism | Best For | Write Amplification | Read Amplification | Space Amplification |
|---|---|---|---|---|---|
| Size-Tiered (STCS) | Merges SSTables of similar size into larger ones. | Write-heavy, append-only workloads (e.g., logs, events). | Low | High | High |
| Leveled (LCS) | Organizes SSTables into levels. Data is compacted into progressively larger, non-overlapping levels. | Read-heavy or update/delete-heavy workloads. | High | Low | Low |
| Time-Window (TWCS) | Groups SSTables into time-based windows and only compacts within a window. | Time-series data with a TTL. | Low | Low (within a window) | Low (for expired data) |
- STCS is the classic write-optimized strategy. It minimizes WA by rewriting data infrequently, but this comes at the cost of potentially having many overlapping SSTables, which hurts read performance.
- LCS is the read-optimized counterpart. It aggressively compacts data to ensure a partition exists in a minimal number of SSTables, providing predictable read latency. This constant rewriting results in high WA.
- TWCS is a specialized form of STCS for time-series data. By isolating data into time windows (e.g., one day), it ensures that queries for a specific time range only hit a small set of SSTables. Once a window expires, its SSTables can be dropped entirely without a costly merge process, making it extremely efficient for TTL'd data.
3. The Modern Approach: Unified Compaction Strategy (UCS)
Recognizing that most workloads are not purely read- or write-heavy, and that switching between the legacy strategies requires a disruptive full-data recompaction, Cassandra introduced the Unified Compaction Strategy (UCS). It is the recommended strategy for most workloads in modern Cassandra versions.
UCS provides a single, flexible framework that can be configured to behave like STCS, LCS, or TWCS, and can even be adjusted on-the-fly without a full recompaction.
Unified Compaction Strategy (UCS) - Apache Cassandra
Let's explore UCS. It's a more sophisticated strategy designed to give operators fine-grained control over the RA/WA trade-off. We'll start with an overview and then look at concrete configurations.
Please read the following sections from the UCS documentation: The introduction under the main heading to understand its purpose. The "Read and write amplification" section to see how the scaling_parameter works. The "Use Case Specific Configurations" table. This is the most practical part, showing how to tune UCS for different workloads.
Key Concepts of UCS
-
Scaling Parameter (
w): This single parameter is the primary knob for tuning the RA/WA balance.- Tiered-like (e.g.,
T4,T8): A positivewvalue (T_fwheref = 2+w) mimics STCS. It prioritizes low write amplification, making it suitable for write-heavy workloads. A higher fanout factorfmeans fewer, larger compactions. - Leveled-like (e.g.,
L10): A negativewvalue (L_fwheref = 2-w) mimics LCS. It prioritizes low read amplification, ideal for read-heavy workloads.
- Tiered-like (e.g.,
-
Sharding: This is a major innovation in UCS. Instead of creating single, monolithic SSTables, UCS can split the token range into multiple shards and compact them in parallel. This has two significant benefits:
- Parallelism: Compactions can run concurrently on different shards, increasing overall compaction throughput.
- SSTable Size Control: It keeps individual SSTable sizes manageable (controlled by
target_sstable_size), which is crucial for operations like streaming, repair, and avoiding long-running I/O-intensive compactions on high-density nodes.
Configuring UCS for Your Workload
The documentation provides excellent starting points. Here’s how you would apply them:
-
For a Read-Heavy Workload (e.g., user profile store):
You want to minimize read amplification. You would configure UCS to behave like LCS.ALTER TABLE my_keyspace.my_table WITH compaction = { 'class': 'UnifiedCompactionStrategy', 'scaling_parameters': 'L10', 'target_sstable_size': '256MiB' };This sets a leveled-like behavior with a fanout factor of 10 and aims for smaller SSTables to keep reads fast.
-
For a Write-Heavy Workload (e.g., event logging):
You want to minimize write amplification. You would configure UCS to behave like STCS.ALTER TABLE my_keyspace.my_table WITH compaction = { 'class': 'UnifiedCompactionStrategy', 'scaling_parameters': 'T4', 'target_sstable_size': '1GiB' };This sets a tiered-like behavior with a fanout factor of 4 and allows for larger SSTables, reducing the frequency of compactions.
-
For a Time-Series Workload (e.g., metrics):
You want the benefits of TWCS. UCS can be configured for this by using a tiered scaling parameter and ensuring TTLs are set on the data.-- Assuming TTL is set on inserted data ALTER TABLE my_keyspace.my_table WITH compaction = { 'class': 'UnifiedCompactionStrategy', 'scaling_parameters': 'T8', 'expired_sstable_check_frequency_seconds': '600' };The tiered approach groups data written around the same time, and the frequent check for expired SSTables helps reclaim space efficiently, similar to how TWCS drops old time windows.
4. Practical Management and Monitoring
Configuring the strategy is only the first step. As an operator, you also need to monitor and manage the compaction process.
Compaction | Apache Cassandra Documentation
The nodetool utility is your primary interface for interacting with compaction on a live node. Let's review the essential commands.
Please review the list under "Compaction nodetool commands". Pay special attention to compactionstats and setcompactionthroughput.
Key commands include:
nodetool compactionstats: Provides a real-time view of pending compactions, active compactions, and their progress. This is your go-to command to check if compaction is keeping up with the write load. A constantly growing number of pending compactions is a sign of trouble.nodetool setcompactionthroughput -- <value_in_mb>: Allows you to throttle the I/O used by compaction. This is a critical safety valve. If compaction is impacting your application's latency, you can lower this value. Conversely, during off-peak hours, you might increase it to catch up on pending compactions.nodetool disableautocompaction/enableautocompaction: Allows you to temporarily pause and resume compaction on a node, which can be useful during critical maintenance or troubleshooting.
For advanced tuning and experimentation, you can also change compaction parameters dynamically via JMX without restarting the node, as mentioned in the documentation.
Conclusion
We've seen that compaction is a complex but essential process in Cassandra, directly governing the trade-off between read and write performance.
Key Takeaways:
- Compaction strategies are chosen to balance Read Amplification, Write Amplification, and Space Amplification based on the application's workload.
- The legacy strategies—STCS (write-heavy), LCS (read-heavy), and TWCS (time-series)—are optimized for specific use cases.
- The modern Unified Compaction Strategy (UCS) is the recommended approach, offering a single, flexible framework that can be tuned to emulate the legacy strategies and provides significant performance benefits through sharding and parallel execution.
- Configuring UCS involves setting the
scaling_parameters(Lfor leveled,Tfor tiered) and other options liketarget_sstable_sizeto match your workload. - Active monitoring with
nodetool compactionstatsand management withnodetool setcompactionthroughputare crucial operational practices.
Preview of the Next Lesson:
This concludes our deep dive into Cassandra. We've covered its distributed architecture, consistency model, and internal storage engine. We will now shift our focus to another critical component of modern distributed systems: in-memory data stores. In the next lesson, we will begin our study of Redis, exploring its persistence models (RDB and AOF) and how they balance durability and performance.