SSTABLE - Immutable Disk Storage
💿 Master Cassandra's on-disk format: file structure, immutability, compaction strategies, bloom filters, and read optimization!
📖 Spotify: 500 Billion SSTables, 100 PB Storage
Spotify's challenge: 400M+ users streaming billions of songs. Every play = write to listening history. Scale: 100 petabytes of SSTable data! The problem: How to store and query billions of listening events efficiently? Traditional databases = random I/O nightmare (10-50ms per read). The solution: SSTables - immutable sorted files. How it works: Memtable flushes to SSTable → 128MB sorted file written sequentially to disk. Immutable once written → Never modified, only merged via compaction. Bloom filters → 99% accurate test "is key in this SSTable?" prevents wasted disk reads. Compression → LZ4 achieves 3-5x compression ratio. Result: P99 read latency: 15ms ✓ (10ms disk seek + 5ms decompression). 500 billion SSTables managed → Compaction keeps it under control. 80% disk space saved → Compression + compaction remove old versions. Reads scale horizontally → Add more nodes, more SSTables handled. Spotify Engineering: "SSTables are the foundation of our scale. Immutability makes everything simple - no locks, no corruption, just sequential writes and smart reads. Bloom filters save us billions of disk I/O operations daily. Without SSTables, Spotify couldn't exist!"
💿 What is an SSTable?
SSTable (Sorted String Table) is Cassandra's immutable on-disk storage format - the permanent home for your data!
Definition
What: Immutable sorted file on disk
Type: Permanent storage (part of LSM tree)
Purpose: Persist data durably
Format: Multiple component files
Sorted: By partition key → clustering key
Immutable: Never modified after creation!
Creation Process
Trigger: Memtable reaches threshold (64-128MB)
Freeze: Memtable becomes read-only
Flush: Write to disk sequentially
Duration: 2-5 seconds typical
Size: Same as memtable (before compression)
Result: New SSTable on disk!
Immutability
Written once: Never modified
Updates: Create new version in new SSTable
Deletes: Tombstone markers
Benefits: No locks, no corruption
Cleanup: Via compaction
Simple: Append-only world!
LSM Tree Level
Level 0: Memtable (memory)
Level 1+: SSTables (disk)
Multiple levels: Organized by size/age
Compaction: Merges levels
Reads: Check all levels
Writes: Only to Level 0 (memtable)
Storage Characteristics
Location: /var/lib/cassandra/data
Per-table: Directory per keyspace.table
Size: 100MB-2GB typical per SSTable
Count: 10-1000 SSTables per table
Compression: LZ4, Snappy, or Deflate
Monitoring: Track count and total size
Performance Trade-offs
Writes: Very fast (sequential flush)
Reads: Slower (must check multiple SSTables)
Space: Can have duplicates (old versions)
Compaction: Cleans up but uses I/O
Optimization: Fewer SSTables = faster reads
Balance: Write speed vs read speed
💡 Why Immutability is Powerful
Immutability solves SO many problems:
No locks needed: Multiple readers can access simultaneously (file never changes!)
No corruption risk: Can't partially update a file that's never updated
Simple crash recovery: SSTables either exist completely or not at all
Easy replication: Just copy files (they never change)
Predictable performance: Read behavior is constant
This is why modern databases (Cassandra, RocksDB, BigTable) all use immutable SSTables!
📚 Core Concepts Explained
Before diving deeper, let's clarify fundamental concepts you need to understand SSTables!
Partition Key & Clustering Key
Partition Key: Primary identifier (WHERE the data lives)
Example: user_id in users table
Purpose: Determines which node stores data
Clustering Key: Secondary sort order (HOW rows are ordered)
Example: timestamp in events table
Purpose: Sorts rows within partition
Together: (user_id, timestamp) = unique row
LSM Tree (Log-Structured Merge Tree)
Concept: Multi-level data structure
Level 0 (L0): Memtable (in memory)
Level 1+ (L1-LN): SSTables (on disk)
Flow: Writes → L0 → flush → L1 → compact → L2...
Why "Log-Structured": Append-only writes
Why "Merge": Periodic compaction merges levels
Key idea: Fast writes (L0), organized reads (L1+)
Tombstones (Deletion Markers)
What: Special marker that says "this is deleted"
Why not actually delete: SSTables are immutable!
How it works: Write tombstone with timestamp
Reads: Skip rows with tombstones
Cleanup: Removed during compaction
Example: DELETE user WHERE id=123
→ Writes tombstone for user 123
→ Actual data removed later by compaction
Sequential vs Random I/O
Sequential I/O: Reading/writing consecutive blocks
Speed: 100-300 MB/s (disk head stays in place!)
Example: Writing memtable to SSTable
Random I/O: Jumping around to different locations
Speed: 1-2 MB/s (disk head moves constantly!)
Example: Reading 10 random keys
Difference: Sequential = 100-300x faster!
Why: Disk seek time (5-10ms) dominates
Amplification Factors
Write Amplification: Total bytes written ÷ user writes
Example: Write 1GB → rewritten 10GB during compaction = 10x
Read Amplification: SSTables checked per read
Example: Check 5 SSTables to find one key = 5x
Space Amplification: Actual space ÷ logical data
Example: 1GB data, 2GB disk used = 2x
Goal: Minimize all three (trade-offs!)
Index Types: Sparse vs Dense
Dense Index: Entry for EVERY key
Size: Huge (can't fit in memory)
Lookup: Direct (fast but impractical)
Sparse Index: Entry for every Nth key
Size: Much smaller (fits in memory!)
Lookup: Find nearest, then scan
SSTable uses: Sparse (1 entry per 128KB)
Why: 99% smaller, still fast enough
False Positive Explained
Concept: Test says "yes" when answer is "no"
Example: Pregnancy test shows positive when not pregnant
In Bloom Filters:
1% false positive = 1 out of 100 times says "key might be here" when it's NOT
True Negative: Says "not here" = ALWAYS correct ✓
False Positive: Says "maybe here" = Sometimes wrong ✗
Impact: 1% wasted disk reads (acceptable!)
Compression Ratio Explained
"3x compression" means: Data becomes 3x SMALLER
Example:
• Original: 300MB
• Compressed: 100MB
• Ratio: 3x (or 3:1)
LZ4: Typically 2-3x smaller
Deflate: Typically 4-5x smaller
Benefit: 3x compression = read 3x less data from disk!
TTL (Time To Live)
Concept: Automatic data expiration
Example: Session data expires after 24 hours
How it works:
• Set TTL when writing: TTL = 86400 (seconds)
• Cassandra tracks creation timestamp
• After 24h: Data automatically deleted
Use cases: Logs, sessions, temporary data
TWCS: Optimized for TTL workloads
Hash Function Basics
What: Function that converts data → number
Example: hash("alice") = 12345
Properties:
• Same input = same output (deterministic)
• Different input = different output (usually)
• Fast to compute
In Bloom Filters:
Use 3-5 different hash functions to set bits
Why multiple: Reduces collision chance
Binary Search Explained
Concept: Efficient search in sorted data
How it works:
1. Start in middle
2. Is key less than or greater than middle?
3. Discard half, repeat
Example: Find 67 in [10, 20, 30, 40, 50, 60, 70, 80, 90]
• Check 50 (middle): 67 > 50
• Check 70: 67 < 70
• Check 60: 67 > 60
• Found between 60-70!
Speed: O(log n) - very fast!
CRC32 Checksum
What: Error detection code
Purpose: Detect data corruption
How it works:
1. Calculate checksum when writing (e.g., 3FA2B8C9)
2. Store checksum with data
3. When reading: recalculate checksum
4. If different → data corrupted!
CRC32: 32-bit checksum (4 bytes)
Accuracy: Detects 99.9999% of errors
🎓 Quick Reference Summary
Remember these key concepts:
Data Organization: Partition key (which node) + Clustering key (sort order)
Storage Structure: LSM Tree = Memtable (L0) → SSTables (L1+)
Deletion: Tombstones (markers) not actual deletion
I/O Speed: Sequential (fast) vs Random (slow) - 100x difference!
Amplification: Write/Read/Space - how much overhead
Index: Sparse (every Nth key) saves memory
Bloom Filter: 1% false positive = 99% accurate "not here"
Compression: 3x = 3 times SMALLER
TTL: Auto-expiration timer
Hash: Convert data → number
Binary Search: Divide-and-conquer = fast lookup
CRC32: Detects corruption
Now you're ready to understand SSTables deeply! 🚀
📁 File Structure: SSTable Components
An SSTable isn't just one file - it's a collection of component files working together!
Data.db
Purpose: Actual data (columns, values)
Format: Binary sorted by partition + clustering key
Compressed: Yes (LZ4/Snappy)
Size: Largest component (GB typical)
Read: Via Index.db for position
Sequential: Sorted = efficient scans
Index.db
Purpose: Maps partition key → Data.db offset
Format: Sparse index (every N partitions)
Typical: One entry per 128KB of data
Size: ~1% of Data.db
Usage: Find data position without full scan
Critical: Enables fast lookups!
Filter.db (Bloom Filter)
Purpose: Probabilistic test "is key here?"
Accuracy: 99% false positive rate configurable
Size: ~10MB per 1GB of data
Loaded: In memory (always)
Saves: Billions of wasted disk reads!
Critical: Read performance guardian
Summary.db
Purpose: Index of the index (sparse)
Format: Every Nth partition key sample
Loaded: In memory (always)
Size: Very small (~1MB typical)
Usage: Narrow Index.db search
Optimization: Reduce Index.db lookups
Statistics.db
Purpose: Metadata about SSTable
Contains: Min/max timestamps, partition count
Compression: Compression ratio info
Tombstones: Deletion count
Usage: Compaction decisions
Size: Few KB only
CompressionInfo.db
Purpose: Compression metadata
Contains: Chunk boundaries, sizes
Chunk size: 64KB default
Usage: Decompress specific chunks
Efficiency: Don't decompress entire file
Random access: Into compressed data
TOC.txt
Purpose: Table of contents
Lists: All component files
Format: Plain text
Usage: Verify SSTable completeness
Integrity: Check all components present
Startup: Load SSTables from TOC
Digest.crc32
Purpose: Checksum for Data.db
Algorithm: CRC32
Validation: Detect corruption
Read time: Verify on access
Scrub: nodetool scrub checks all
Reliability: Data integrity guarantee
📁 Example SSTable Files
Filename format: {version}-{generation}-{component}
/var/lib/cassandra/data/mykeyspace/users-abc123/
├── md-1-big-Data.db (2.1 GB - actual data)
├── md-1-big-Index.db (21 MB - partition index)
├── md-1-big-Filter.db (18 MB - bloom filter)
├── md-1-big-Summary.db (450 KB - index summary)
├── md-1-big-Statistics.db (12 KB - metadata)
├── md-1-big-CompressionInfo.db (8 KB - compression map)
├── md-1-big-TOC.txt (100 bytes - file list)
└── md-1-big-Digest.crc32 (32 bytes - checksum)
Total: 2.15 GB (2.1 GB data + 50 MB overhead = 2.4% overhead)
🔒 Immutability: The Power of Write-Once
Immutability is the secret sauce that makes SSTables so powerful!
No In-Place Updates
Traditional DB: Update record in place
SSTable: Write new version to new SSTable
Old version: Remains until compaction
Reads: Return newest timestamp
Cleanup: Compaction removes old versions
Benefit: No file locking needed!
Lock-Free Reads
Multiple readers: Can access simultaneously
No coordination: File never changes!
Concurrent: Unlimited read parallelism
No contention: No lock waits
Predictable: Constant performance
Throughput: Scales with CPU cores
Crash Safety
Complete or not: SSTable fully written or doesn't exist
No partial: Can't corrupt existing SSTables
Atomic: Rename operation creates SSTable
Recovery: Just delete incomplete ones
Simple: No complex recovery logic
Reliable: Never corrupt existing data
Easy Replication
File copy: Just copy SSTable files
No locking: Copy while reads continue
Streaming: Transfer to new node
Checksum: Verify after transfer
Simple: Standard file operations
Repair: Replace damaged files easily
Compaction Friendly
Background merge: Happens async
Old SSTables: Keep serving reads
New SSTable: Created separately
Atomic swap: Switch to new, delete old
No downtime: Zero impact on reads
Clean: Simple file operations
Time-Travel Reads
Versioning: Old versions remain
Timestamp: Each mutation has time
Point-in-time: Read as-of timestamp
History: Before compaction removes
Debugging: See what data was when
Compliance: Audit trail possible
🎯 Netflix: Zero Corruption in 5 Years
Netflix's reliability achievement:
The Challenge:
• 1,000+ Cassandra nodes
• Billions of SSTable files created over 5 years
• Hardware failures, power outages, network issues
• Kernel crashes, disk failures
The Result:
• Zero SSTable corruption incidents ✓
• Every crash = clean recovery
• Incomplete SSTables = simply deleted
• Existing SSTables = never damaged
Why:
Immutability means crashes can't corrupt existing data. Either the new SSTable is complete (atomic rename succeeded) or it doesn't exist (temp file deleted). There's no middle ground where an existing SSTable gets partially updated and corrupted.
Netflix Engineering: "In 5 years of production Cassandra, we've never seen a single case of SSTable corruption from crashes. Immutability makes corruption impossible. This is worth any trade-off in extra compaction work!"
⚖️ The Trade-off: Space for Safety
Immutability isn't free:
Cost: Multiple versions of same data (old + new). Updates don't replace, they add. Deletes create tombstones (more data!). Requires compaction to reclaim space.
Benefit: Zero corruption risk. Lock-free reads. Simple crash recovery. Easy replication. Predictable performance.
Verdict: Worth it! Disk is cheap, corruption is expensive. Compaction runs in background and cleans up. The safety and simplicity win!
📖 Read Path: Finding Data in SSTables
Reading from SSTables is carefully optimized to minimize disk I/O!
1. Check Memtable First
Priority: Check in-memory memtable
Hit: Return immediately (sub-ms!)
Miss: Continue to SSTables
Recent data: 80-95% hit rate
Fast: No disk I/O
Next: Check SSTables newest → oldest
2. Bloom Filter Test
Check: Is partition key in this SSTable?
In memory: Bloom filter always loaded
False positive: 1% chance says "yes" when not there
True negative: 100% accurate "no"
Saves: 99% of wasted disk reads!
Critical: Read performance guardian
3. Key Cache Lookup
Check: Do we know position already?
Cache: Maps partition key → Data.db offset
Hit: Skip Index.db read!
Size: Configurable (200MB-2GB typical)
Hit rate: 80-90% for hot data
Saves: One disk read per lookup
4. Index.db Lookup
If cache miss: Read Index.db
Binary search: Find partition key
Get offset: Position in Data.db
Latency: ~5-10ms disk read
Cache: Store for next time
Optimization: Index.db in OS page cache
5. Data.db Read
Seek: Position from Index.db
Read chunk: 64KB compressed block
Decompress: LZ4 decompression
Parse: Extract row data
Latency: 10-20ms total
OS cache: Hot data cached
6. Merge Results
Multiple SSTables: Check all (newest first)
Merge: Combine by timestamp
Newest wins: Latest timestamp returned
Tombstones: Mark deleted data
Result: Complete row view
Return: To client
⚡ Bloom Filter: The Unsung Hero
Without bloom filters: Would need to check Index.db for every SSTable (~10ms each). 10 SSTables = 100ms wasted on disk I/O for keys that aren't even there!
With bloom filters: 1ms memory check says "definitely not here" for 99% of SSTables. Result: 100ms → 1ms (100x faster!)
Real impact: Spotify saves an estimated 10 billion disk I/O operations per day thanks to bloom filters. This is the difference between 15ms and 500ms+ P99 read latency!
🌸 Bloom Filters: Probabilistic Magic
Bloom filters are the secret weapon that makes SSTable reads fast!
How Bloom Filters Work
Data structure: Bit array + hash functions
Insert: Hash key → set multiple bits to 1
Test: Hash key → check if all bits are 1
All 1s: "Maybe present" (could be false positive)
Any 0: "Definitely not present" (100% accurate!)
Trade-off: Space vs accuracy
False Positive Rate
Default: 1% false positive rate
Meaning: 1% chance says "maybe" when not there
Configurable: 0.01% - 10%
Lower FP: Larger bloom filter
Higher FP: Smaller but more disk I/O
Sweet spot: 1% for most workloads
Memory Usage
Formula: ~10 bytes per key for 1% FP rate
1M keys: ~10MB bloom filter
100M keys: ~1GB bloom filter
Loaded: Always in memory
Critical: Must fit in RAM!
Monitor: Total bloom filter size
Performance Impact
Check time: <1ms (memory lookup)
Hash operations: 3-5 hashes typical
Saves: 99% of unnecessary disk reads!
10 SSTables: 9.9 skipped on average
Result: 100ms → 1ms for misses
Worth it: Massive performance win
Configuration
bloom_filter_fp_chance: In table schema
Default: 0.01 (1% false positive)
Read-heavy: 0.001 (0.1% - larger filter)
Write-heavy: 0.1 (10% - smaller filter)
Disable: 1.0 (not recommended!)
Per-table: Can tune differently
Tuning Trade-offs
Lower FP (0.1%):
✅ Fewer wasted disk reads
❌ 10x larger bloom filters
Higher FP (10%):
✅ Smaller memory usage
❌ More unnecessary disk I/O
Default 1%: Best balance!
📊 Apple: 10 Billion Daily Bloom Filter Checks
Apple's iCloud scale:
The Numbers:
• 1 billion+ active users
• 10 billion SSTable reads per day
• Average 20 SSTables per read query
• = 200 billion potential disk reads!
Bloom Filter Impact:
• 99% filter efficiency
• 200B potential reads → 2B actual reads
• 198 billion disk I/O operations saved per day!
• Average 10ms per avoided read = 22,916 days saved!
• Memory cost: 500GB bloom filters across cluster
Configuration:
• bloom_filter_fp_chance: 0.005 (0.5%)
• Why lower: Read-heavy workload
• Cost: 2x bloom filter size vs 1% FP
• Worth it: Saves billions more I/O ops
Apple Engineering: "Bloom filters are non-negotiable. The memory cost is tiny compared to the disk I/O savings. At our scale, without bloom filters, we'd need 100x more storage servers just to handle the I/O load!"
🧮 Bloom Filter Math
Memory calculation for 1% false positive rate:
Bits per key = -log(FP) / (ln(2))²
= -log(0.01) / 0.48 ≈ 9.6 bits per key
Example: 100 million keys
Memory = 100M × 9.6 bits = 960M bits = 120 MB
Lower FP (0.1%):
Memory = 100M × 14.4 bits = 180 MB (50% more)
Higher FP (10%):
Memory = 100M × 4.8 bits = 60 MB (50% less)
Rule of thumb: ~10 bytes per key for 1% FP rate
🔄 Compaction: Cleanup and Optimization
Compaction merges SSTables to reclaim space and improve read performance!
Why Compaction?
Problem: SSTables accumulate over time
Duplicates: Old versions of updated rows
Tombstones: Deleted data markers
Reads slow: More SSTables to check
Space waste: Old data not removed
Solution: Merge and clean!
Compaction Process
Select: Pick SSTables to merge
Read: Stream through all selected
Merge: Combine by key, newest wins
Remove: Old versions and tombstones
Write: New compacted SSTable(s)
Delete: Old SSTables after complete
STCS: Size-Tiered
Default strategy for most workloads
How: Merge SSTables of similar size
Trigger: 4+ SSTables of same tier
Pros: Great for write-heavy, simple
Cons: Read amp high, space amp 2x
Best for: Time-series, logs, IoT
LCS: Leveled
Read-optimized strategy
How: Organize into levels by size
L0: New SSTables
L1-LN: Non-overlapping ranges
Pros: Minimal read amp, 90% less space
Cons: 10x more write I/O
Best for: Read-heavy workloads
TWCS: Time-Window
Time-series optimized
How: Compact by time window
Windows: Hour, day, week buckets
Expire: Drop entire SSTables (fast!)
Pros: Efficient TTL, predictable
Cons: Only for time-series
Best for: Metrics, events, logs with TTL
Strategy Comparison
STCS: Write 1GB → Read 1-10 SSTables
LCS: Write 1GB → Read 1-2 SSTables ✓
TWCS: Write 1GB → Read 1 SSTable per window
Choose based on: Read vs write ratio
Most common: STCS (default)
⚠️ Compaction Costs
Compaction isn't free:
I/O cost: Read old SSTables + write new ones. Can consume 50% of disk bandwidth!
CPU cost: Decompression + merge + compression. Can use 20-30% CPU.
Temporary space: Need 2x space during compaction (old + new).
But worth it: Faster reads, less space long-term, remove tombstones. Run during off-peak hours if possible!
📖 Concrete Example: Same Workload, Different Strategies
Scenario: User activity tracking table
• 1 million new events per day
• Each event = 1KB
• Query pattern: Recent events + user history
• Retention: 90 days
Strategy 1: STCS (Size-Tiered)
Day 1: 1GB SSTable created
Day 2: Another 1GB SSTable
Day 3: Another 1GB SSTable
Day 4: Trigger! 4 SSTables → compact into 4GB SSTable
After 90 days: ~50 SSTables of various sizes
Read query: Must check 20-30 SSTables
Disk space: ~180GB (2x amplification)
Best for: Heavy writes, analytics queries OK
Strategy 2: LCS (Leveled)
Levels organized:
• L0: 1GB daily SSTables
• L1: 10GB (10 non-overlapping 1GB files)
• L2: 100GB (100 non-overlapping 1GB files)
After 90 days: ~10 SSTables per read
Read query: Check 1-2 SSTables only! ✓
Disk space: ~100GB (1.1x amplification) ✓
Write amplification: 10x (rewrites during level promotion) ❌
Best for: Read latency critical
Strategy 3: TWCS (Time-Window)
Windows: 1 day per SSTable
Day 1: SSTable_2024-01-01 (1GB)
Day 2: SSTable_2024-01-02 (1GB)
Day 90: SSTable_2024-03-30 (1GB)
Day 91: SSTable_2024-01-01 deleted! (TTL expired) ✓
Query "last 7 days": Check exactly 7 SSTables
Expiration: 1-second file delete (vs hours of compaction!) ✓
Disk space: ~135GB (1.5x amplification)
Best for: Time-series with TTL
Key Takeaway:
Same data, same writes - dramatically different performance based on strategy!
• STCS: Fast writes, slower reads (20-30 SSTables)
• LCS: Slower writes, fast reads (1-2 SSTables)
• TWCS: Best for time-series, instant TTL cleanup
📦 Compression: Save Space, Speed Reads
Compression reduces storage by 3-5x with minimal CPU cost!
LZ4 (Default)
Speed: Extremely fast (GB/s)
Ratio: 2-3x compression
CPU: Very low overhead
Decompress: 500MB/s+ single thread
Best for: Most workloads
Default: Recommended!
Snappy
Speed: Very fast
Ratio: 2-3x (similar to LZ4)
CPU: Low overhead
Decompress: 300-400MB/s
Mature: Google-developed
Use: Alternative to LZ4
Deflate (Zlib)
Speed: Slower compression
Ratio: 4-5x (best ratio!)
CPU: Higher overhead
Decompress: 100-200MB/s
Best for: Cold data, archival
Trade-off: Space vs CPU
Chunk Size
Default: 64KB chunks
Why chunk: Random access into file
Read: Decompress only needed chunk
Smaller (16KB): Less decompression work
Larger (256KB): Better ratio
Sweet spot: 64KB default works well
Performance Impact
Disk I/O: 3x less data to read ✓
CPU: Decompression cost (2-5ms)
Net effect: Usually faster!
Why: Disk slower than CPU
SSD: LZ4 still wins
HDD: Any compression wins big
Choosing Algorithm
Hot data: LZ4 (fast reads)
Warm data: LZ4 or Snappy
Cold data: Deflate (save space)
Default LZ4: 99% of cases
Change: Per-table configuration
Benchmark: Test your data!
📈 Compression Math
Example: 1GB SSTable with LZ4 (3x compression)
Storage saved: 1GB → 333MB (667MB saved)
Disk I/O: Read 333MB instead of 1GB (3x faster!)
Decompression: 333MB @ 500MB/s = 0.67 seconds
Uncompressed read: 1GB @ 150MB/s = 6.7 seconds
Result: 0.67s + minimal CPU vs 6.7s = 10x faster!
Even with decompression overhead, compressed reads are much faster because disk I/O is the bottleneck!
🚀 Performance Characteristics
Understanding SSTable performance helps with optimization and capacity planning!
Write Performance
Memtable flush: Sequential write
Speed: 100-300 MB/s (SSD)
Duration: 2-5 seconds for 128MB
Non-blocking: New memtable created
Bottleneck: Disk bandwidth
Optimization: Larger memtables
Read Performance
Best case: Bloom filter says no (1ms)
Cache hit: Key cache + OS page cache (2-5ms)
Cold read: Index + data (10-50ms)
Multiple SSTables: Check each (20-100ms)
Optimization: Compaction reduces count
Critical: Fewer SSTables = faster
Read Amplification
Definition: SSTables checked per read
STCS: 5-20 SSTables typical
LCS: 1-2 SSTables (optimized!)
TWCS: 1-3 per time window
Impact: Linear on read latency
Mitigation: Choose right strategy
Space Amplification
Definition: Actual vs logical space
STCS: 2x space (old + new versions)
LCS: 1.1x space (minimal waste!)
TWCS: 1.5x space
Compaction: Reclaims space
Monitoring: Watch disk usage
Write Amplification
Definition: Total bytes written vs user writes
STCS: 5-10x (rewrite during compaction)
LCS: 10-30x (more compaction!)
TWCS: 2-5x (minimal rewriting)
Trade-off: Read vs write optimization
SSD wear: LCS impacts lifespan
Optimization Tips
Reads slow: Use LCS, tune bloom filters
Writes slow: Use STCS, larger memtables
Disk full: Run compaction, use compression
High latency: Reduce SSTable count
Monitor: SSTable count per read
Test: Benchmark your workload!
🏢 Real Company SSTable Strategies
📊 Discord: LCS for 1 Trillion Messages
Use case: Message storage (1 trillion+ messages)
Challenge: Read-heavy (users scroll chat history)
Strategy: Leveled Compaction (LCS)
• Why: Minimize read amplification
• sstable_size_in_mb: 160MB (smaller than default)
• Levels: 0-6 typical depth
• Non-overlapping ranges = predictable reads
Results:
• Read amplification: 1-2 SSTables per read (vs 10-20 with STCS!)
• P99 read latency: 15ms (consistent)
• Space usage: 10% more efficient than STCS
• Trade-off: 10x more compaction I/O (acceptable!)
Configuration:
• compression: LZ4 (fast decompression)
• bloom_filter_fp_chance: 0.01 (1%)
• key_cache_size: 1GB per node
Discord Engineering: "LCS transformed our read performance. Scroll lag disappeared. The extra compaction I/O is worth it - users care about read latency, not backend I/O!"
⏰ Uber: TWCS for Time-Series Trips
Use case: Trip history (location pings every 5 seconds)
Pattern: Write recent, read recent, expire old (90 days TTL)
Strategy: Time-Window Compaction (TWCS)
• compaction_window_unit: DAYS
• compaction_window_size: 1 (1-day windows)
• TTL: 90 days
• Why: Efficient expiration + predictable queries
How it works:
• Day 1 data → SSTable_day1
• Day 2 data → SSTable_day2
• Query "last week": Check 7 SSTables
• TTL expired: Drop entire SSTable (instant!)
Results:
• Expiration: Delete entire file (1 second vs hours of compaction)
• Space reclaimed: Immediately
• Read performance: Predictable (know which SSTables to check)
• Write amplification: 2x (minimal rewriting)
Uber Engineering: "TWCS is perfect for time-series. Dropping expired data used to take hours of compaction I/O. Now it's instant - just delete the file. This saved us millions in storage costs!"
🎮 Epic Games: STCS for Write-Heavy Telemetry
Use case: Fortnite game telemetry (100M events/hour)
Pattern: Massive writes, rare reads (analytics only)
Strategy: Size-Tiered Compaction (STCS - Default)
• Why: Optimized for write throughput
• min_threshold: 4 (default)
• max_threshold: 32 (default)
• No tuning needed!
Why STCS wins:
• Minimal write amplification (5-10x vs 30x LCS)
• Fast memtable flushes (no blocking)
• Compaction in background (doesn't slow writes)
• Reads are rare (analytics batch jobs)
Configuration:
• memtable_heap_space: 256MB (2x default)
• compression: LZ4 (4x ratio on telemetry)
• bloom_filter_fp_chance: 0.1 (10% - don't care about reads)
Results:
• Write throughput: 500k events/sec per node ✓
• Compaction overhead: 20% disk I/O
• Read latency: 50-100ms (acceptable for batch)
• Storage: 2x amplification (worth it for write speed)
Epic Engineering: "Don't over-optimize! STCS default settings handle our 100M events/hour perfectly. We tried LCS and writes slowed 3x. For write-heavy workloads, STCS is king!"
✅ Best Practices: SSTable Optimization
1. Choose Right Compaction Strategy
Read-heavy: LCS (leveled)
Write-heavy: STCS (size-tiered)
Time-series + TTL: TWCS (time-window)
Test both: Benchmark your workload
Monitor: Read amp, write amp, space amp
Don't assume: One size doesn't fit all
2. Monitor SSTable Count
Healthy: <10 SSTables per table
Warning: 10-50 SSTables
Critical: >50 SSTables (slow reads!)
Check: nodetool tablestats
Fix: Force compaction if stuck
Alert: Set up monitoring
3. Tune Bloom Filters
Default 1%: Good for most
Read-heavy: 0.1% (larger filter, fewer false positives)
Write-heavy: 1-10% (smaller filter)
Monitor: Bloom filter size in memory
Trade-off: Memory vs disk I/O
Per-table: Tune differently
4. Use Compression
Always enable: Default LZ4
Hot data: LZ4 (fast)
Cold data: Deflate (better ratio)
Chunk size: 64KB default works well
Savings: 3-5x disk space
Benefit: Faster reads (less I/O)
5. Schedule Compaction Wisely
Background: Runs automatically
I/O intensive: Can slow queries
Off-peak: Run major compactions at night
Throttle: compaction_throughput_mb_per_sec
Monitor: Pending compactions queue
Critical: Don't disable entirely!
6. Provision Adequate Disk Space
Rule: 50% free space minimum
Compaction: Needs 2x temporary space
STCS: 2x permanent amplification
LCS: 1.1x permanent
Alert: At 70% disk usage
Emergency: Add nodes if >80%
🚨 Common Mistakes
- ❌ Wrong compaction strategy: LCS for write-heavy = death spiral
- ❌ Ignoring SSTable count: >50 SSTables = slow reads
- ❌ Disabling bloom filters: Never set fp_chance = 1.0!
- ❌ No compression: Wasting 70% disk space
- ❌ Insufficient disk space: Compaction fails, cluster dies
- ❌ Disabling compaction: SSTables accumulate forever
- ❌ Not monitoring: Problems invisible until too late
💼 Interview Questions & Answers
Complete Answer:
An SSTable (Sorted String Table) is Cassandra's immutable on-disk storage format - a collection of sorted files that permanently store data flushed from memtables.
What is SSTable:
- Collection of files: Data.db, Index.db, Filter.db, Summary.db, Statistics.db, CompressionInfo.db
- Sorted: By partition key → clustering key for efficient lookups
- Immutable: Written once, never modified
- Compressed: Typically 3-5x compression with LZ4
- Part of LSM tree: Memtable (L0) → SSTables (L1+)
- Created by: Memtable flush when reaches threshold
Why Immutability:
1. Lock-free reads:
- Multiple readers can access simultaneously
- File never changes so no coordination needed
- Unlimited read parallelism
- Predictable performance - no lock contention
2. Crash safety:
- SSTables are atomic - either complete or don't exist
- Can't corrupt existing SSTables during crash
- Simple recovery - just delete incomplete files
- No complex recovery logic needed
- Real example: Netflix zero SSTable corruption in 5 years
3. Easy replication:
- Just copy files to other nodes
- Can copy while reads continue (no locking)
- Streaming for new nodes is simple
- Repair = replace damaged files
4. Compaction-friendly:
- Background merge creates new SSTable
- Old SSTables keep serving reads
- Atomic swap to new, delete old
- Zero downtime during compaction
Trade-offs:
- Cost: Updates create new versions (space amplification)
- Deletes create tombstones (more data temporarily)
- Requires compaction to reclaim space
- Benefit: Worth it! Disk cheap, corruption expensive
Key Insight: Immutability trades some disk space for massive simplicity and reliability. No locks, no corruption, simple replication. Modern databases (Cassandra, RocksDB, BigTable) all use this pattern because the benefits far outweigh the costs!
Complete Answer:
Bloom filters are probabilistic data structures that quickly test whether a partition key exists in an SSTable - saving billions of unnecessary disk reads!
How Bloom Filters Work:
Data structure:
- Bit array (millions of bits set to 0 or 1)
- Multiple hash functions (typically 3-5)
- Each key hashes to multiple bit positions
Insert operation:
- Hash partition key with 3-5 different hash functions
- Each hash produces a bit position
- Set all those bit positions to 1
- Example: key "user123" → bits 45, 892, 3421 set to 1
Test operation:
- Hash the query key with same hash functions
- Check if ALL corresponding bits are 1
- All 1s: "Maybe present" (might be false positive)
- Any 0: "Definitely NOT present" (100% accurate!)
False Positive Rate:
- Default: 1% false positive rate
- Meaning: 1% chance says "maybe" when key not actually there
- True negative: 100% accurate (if it says no, definitely no!)
- Configurable: 0.01% - 10% via bloom_filter_fp_chance
Why Critical for Performance:
The problem without bloom filters:
- Reading key "user999" that doesn't exist
- Must check Index.db for ALL 20 SSTables
- 20 SSTables × 10ms disk read = 200ms wasted!
- Most queries access keys not in most SSTables
With bloom filters:
- Check bloom filter in memory (~1ms)
- Says "definitely not here" for 99% of SSTables
- Skip 19 out of 20 SSTables
- 200ms → 1ms (200x faster!)
Real Impact - Spotify Example:
- 10 billion SSTable reads per day
- Average 20 SSTables checked per read
- = 200 billion potential disk reads
- Bloom filters: 99% efficiency
- 200B → 2B actual reads (198 billion saved!)
- 10ms per saved read = 22,916 days of I/O saved daily
Memory Cost:
- Formula: ~10 bytes per key for 1% FP rate
- 100 million keys = 1GB bloom filter
- Must fit in memory (always loaded)
- Trade-off: Small memory cost for massive I/O savings
Tuning:
- Lower FP (0.1%): Larger filter, fewer false positives
- Use for read-heavy workloads
- Higher FP (10%): Smaller filter, more false positives
- Use for write-heavy workloads where reads are rare
Key Insight: Bloom filters are the unsung heroes of Cassandra performance. For minimal memory cost (~1% of data size), they eliminate 99% of wasted disk I/O. This is THE optimization that makes SSTable-based databases viable at scale!
Complete Answer:
Cassandra has three main compaction strategies, each optimized for different workload patterns.
1. STCS (Size-Tiered Compaction Strategy)
How it works:
- Groups SSTables into "tiers" by similar size
- When 4+ SSTables in same tier → merge them
- Creates larger SSTable at next tier
- Tiers: 10MB, 100MB, 1GB, 10GB, etc.
- Simple and write-optimized
Characteristics:
- Read amplification: 5-20 SSTables per read
- Write amplification: 5-10x (moderate rewriting)
- Space amplification: 2x (old + new versions)
- Compaction frequency: Less frequent (efficient)
When to use STCS:
- Write-heavy workloads: Inserts >> reads
- Time-series data (logs, events, metrics)
- IoT sensor data
- Append-only patterns
- Example: Epic Games telemetry (100M events/hour)
- Default: Good starting point for most workloads
2. LCS (Leveled Compaction Strategy)
How it works:
- Organizes SSTables into levels (L0, L1, L2...)
- Each level 10x larger than previous
- L1+: Non-overlapping key ranges
- Promotes SSTables through levels
- Constant compaction happening
Characteristics:
- Read amplification: 1-2 SSTables per read (optimal!)
- Write amplification: 10-30x (lots of rewriting)
- Space amplification: 1.1x (very efficient)
- Compaction frequency: Continuous (high I/O)
When to use LCS:
- Read-heavy workloads: Reads >> writes
- User profiles, session data
- Applications where read latency is critical
- Disk space constrained (90% less waste than STCS)
- Example: Discord messages (1 trillion+ messages)
- Warning: Don't use for write-heavy (death spiral!)
3. TWCS (Time-Window Compaction Strategy)
How it works:
- Groups SSTables by time window (hour, day, week)
- Each window becomes one SSTable
- Old windows expired by deleting entire SSTable
- Designed specifically for time-series + TTL
Characteristics:
- Read amplification: 1-3 SSTables per time range
- Write amplification: 2-5x (minimal)
- Space amplification: 1.5x
- TTL expiration: Instant (delete file!)
When to use TWCS:
- Time-series with TTL: Must have expiration!
- Metrics, monitoring data
- Application logs
- Event streams
- Queries typically by time range
- Example: Uber trip history (90 day TTL)
- Requirement: Don't use without TTL!
Decision Framework:
Choose STCS if:
- Write:Read ratio > 10:1
- Append-only or mostly inserts
- Reads are batch/analytics (latency OK)
- Uncertain (good default)
Choose LCS if:
- Write:Read ratio < 1:10
- Read latency is critical
- Lots of updates/deletes
- Disk space constrained
Choose TWCS if:
- Time-series data
- Has TTL expiration
- Queries by time range
- Recent data accessed most
Key Insight: There's no universal best strategy! STCS favors writes, LCS favors reads, TWCS favors time-series. Choose based on your workload pattern. When in doubt, start with STCS and switch to LCS only if read latency becomes a problem!
Responsive Ad