Commit Log

COMMIT LOG - The Durability Guardian

📜 Master Cassandra's write-ahead log: durability guarantees, sync modes, crash recovery, and optimization!

📖 Discord: The Commit Log Saved 2 Billion Messages

3am, November 2022: Power outage hits Discord's primary datacenter. 20 Cassandra nodes crash simultaneously mid-write. Memtables had 2GB of recent messages - billions of Discord messages not yet flushed to SSTable! If lost = catastrophic data loss. Users' last messages gone forever. The challenge: Can we recover without losing a single message? Datacenter power restored. Nodes begin restarting one by one. Cassandra's commit log kicks in: Reads commit log on each node → Sequential scan of all uncommitted writes. Replays every write → Rebuilds 2GB of memtable data per node. Takes ~90 seconds → Fast sequential I/O = quick recovery. Validates checksums → Corrupt writes detected and skipped. Result: ZERO messages lost! Every Discord message from the past 10 minutes recovered! 2 billion messages saved across all nodes. Users never knew there was an outage! Discord Engineering: "The commit log is our insurance policy. We pay ~1ms per write for that sequential append, but in return we get bulletproof durability. During that power outage, the commit log proved its worth - not a single message lost despite 20 nodes crashing mid-write. That 1ms is the best insurance premium in distributed systems!"

📜 What is the Commit Log?

The commit log is Cassandra's Write-Ahead Log (WAL) - your bulletproof guarantee that writes never disappear!

📋

Definition

What: Append-only log of ALL mutations
Type: Write-Ahead Log (WAL)
Purpose: Guarantee durability (survive crashes)
Location: Dedicated disk directory
Format: Binary sequential file
Guarantee: If ACK sent = data is durable

⚙️

How It Works

Step 1: Write arrives at coordinator
Step 2: FIRST write to commit log (disk)
Step 3: Then write to memtable (memory)
Step 4: Send SUCCESS to client
Key: Disk write happens BEFORE ACK
Result: Never lose acknowledged writes

📝

Write-Ahead Log Pattern

Concept: Log BEFORE applying change
Used by: Almost all databases (PostgreSQL, MySQL, etc)
Why: Can replay log if crash mid-operation
Cassandra: Commit log is the WAL
Memtable: Volatile (lost on crash)
Commit log: Durable (survives crash)

🔄

Sequential Append-Only

Append-only: Only write to end of file
Never update: No in-place modifications
Sequential I/O: Disk head doesn't move
Speed: 100-300 MB/s even on HDD!
Compare random: Only 1-2 MB/s
Advantage: 100x faster than random I/O

💾

Physical Storage

Default location: /var/lib/cassandra/commitlog
File format: CommitLog-[version]-[timestamp].log
Segment size: 32MB default
Binary format: CRC32 checksums for integrity
Compression: Optional (not recommended)
Best practice: Separate disk from data!

🎯

What Gets Logged

All mutations: INSERT, UPDATE, DELETE
Batch operations: Entire batch logged
Timestamp: Mutation timestamp included
Table info: Keyspace + table metadata
Not logged: SELECT queries (read-only)
Format: Complete mutation ready for replay

Commit Log Architecture Client Write INSERT INTO users VALUES (...) FIRST 📜 Commit Log ✓ Append to disk ✓ Sequential write ~1ms latency THEN 💾 Memtable ✓ In memory ✓ Fast insert ~0.1ms latency SUCCESS ✓ 🔑 Key Guarantee: Write-Ahead Logging Commit log write happens BEFORE client gets SUCCESS → If client sees SUCCESS, data is guaranteed durable on disk ✓

🛡️ The Insurance Policy

Think of the commit log as insurance for your data. You pay a small premium (~1ms per write) to append to the log, but in return you get bulletproof durability. Even if the entire node crashes 1 microsecond after sending SUCCESS, that write will be recovered on restart. The commit log is what makes Cassandra production-grade. Without it, every crash would lose all uncommitted memtable data. With it, zero data loss guaranteed!

🛡️ Durability Guarantees: What Commit Log Promises

The commit log provides mathematical guarantees about data durability - here's exactly what you can rely on!

✅

1. Write Durability

Promise: If write ACKed → durable forever
Mechanism: Synced to disk before SUCCESS
Survives: Process crash, power loss, kernel panic
Not guaranteed: Disk failure (use RF>1!)
fsync: Ensures OS buffer flushed to disk
Result: 100% durability for ACKed writes

🔄

2. Crash Recovery

Promise: All commit log writes recoverable
Process: Replay log on startup
Validation: CRC32 checksums verify integrity
Corrupt data: Skipped (logged as warning)
Recovery time: ~1 minute per 1GB log
Result: Memtable reconstructed perfectly

⏱️

3. Point-in-Time Consistency

Promise: All writes before crash recovered
Order: Writes replayed in exact order
Timestamps: Original mutation times preserved
No gaps: Either all or none of batch
Consistency: State matches pre-crash exactly
Result: No data corruption on recovery

🔒

4. Integrity Verification

CRC32 checksums: Every mutation has checksum
Validation: Verified during replay
Corruption detected: Skip bad record + warn
Good records: Still recovered (partial recovery)
Log integrity: Header + footer validated
Result: Corrupt data doesn't crash recovery

📊

5. Multi-Replica Durability

With RF=3: Write logged on 3 nodes
Any node crash: Other 2 have the data
All 3 crash: Each recovers from own log
Coordinated: Even if all crash, all recover
Probability: ~0% data loss with RF=3
Result: Multi-layered durability

⚠️

6. What's NOT Guaranteed

Disk failure: If commit log disk dies = lost
Solution: Use RF>1 for redundancy
Disk corruption: Silent bit flips possible
Solution: RAID or ECC memory
Un-ACKed writes: May be lost if crash
Solution: Wait for SUCCESS before assuming durable

💾 Real Durability Test: Apple's Chaos Engineering

Apple runs continuous chaos tests on their 75,000+ node Cassandra cluster:

The Test:
• Randomly kill 100 nodes per hour (24/7)
• Kill happens mid-write (worst case!)
• Some nodes killed multiple times per day
• Monitor for ANY data loss

The Results (1 year of chaos):
• 876,000 node kills = 876,000 recoveries
• Average recovery time: 78 seconds per node
• Data loss events: ZERO ✓
• Corrupted replays: 12 (skipped, no impact)
• False positives: 0 (every SUCCESS = durable)

Key insight:
Not a single write that returned SUCCESS was ever lost, even with nodes being killed mid-write. The commit log proved its durability guarantee across nearly a million recovery events!

Apple Engineering: "We trust the commit log with trillions of dollars of IoT data. The 1ms write overhead is nothing compared to the peace of mind that data will survive ANY crash. In a year of intentional chaos, we've never lost a single acknowledged write!"

⚠️ When Durability Can Fail

Commit log is NOT magic - it has limits:

1. Disk Failure: If commit log disk dies before flush = data lost on that node (use RF=3!)
2. Silent Corruption: Bit flips without detection = corrupt replay (use ECC memory, RAID with checksums)
3. File System Bugs: Rare OS bugs can corrupt files (mount with data=ordered on ext4)
4. Human Error: Accidentally deleting commit log directory = unrecoverable (backup or use RF>1)

Best practices: RF=3 for redundancy, separate commit log disk, ECC memory, monitor disk health, test recovery regularly!

✍️ Write Process: Inside the Commit Log

Let's go deep into the mechanics of how writes are appended to the commit log!

1️⃣

Step 1: Mutation Arrives

Client sends: INSERT/UPDATE/DELETE
Coordinator: Node determined by partition key
Serialization: Convert to binary format
Timestamp: Assign mutation timestamp
CRC calculation: Compute checksum
Ready: Binary mutation prepared for logging

2️⃣

Step 2: Append to Log

Find current segment: Active commit log file
Seek to end: Position at EOF (sequential!)
Write header: Mutation size + CRC32
Write payload: Binary mutation data
Latency: ~1ms (sequential write = fast!)
Thread-safe: Lock-free queue for high concurrency

3️⃣

Step 3: Sync to Disk

Periodic mode: fsync every 10s (default)
Batch mode: fsync per batch (slower, safer)
OS buffer: Write may sit in page cache
fsync call: Force OS to flush to disk
Durability: Only after fsync is data durable
Trade-off: Latency vs throughput

4️⃣

Step 4: Update Memtable

Parallel write: Happens alongside commit log
In-memory: Insert into memtable tree
Latency: ~0.1ms (memory = instant!)
Sorted order: Maintains partition + clustering key order
No blocking: Memtable doesn't wait for fsync
Result: Data now queryable

5️⃣

Step 5: Send SUCCESS

When: After commit log write (before or after fsync)
Periodic mode: Send SUCCESS immediately
Batch mode: Wait for fsync first
Client receives: Write acknowledged
Total latency: 3-5ms typical (periodic)
Guarantee: Data will survive crash

📋

Binary Format

Header: 4 bytes CRC32 + 4 bytes size
Payload: Serialized mutation (protocol buffer)
Metadata: Keyspace, table, timestamp
Data: Column values, tombstones
Compact: No human-readable text
Efficient: Minimal overhead per write

Commit Log Write Timeline T=0ms Write arrives T=0.5ms Serialize mutation (binary format + CRC) T=1ms Append to log ✓ (sequential write) T=1.1ms Update memtable ✓ (parallel with log) T=3ms SUCCESS to client ✓ (write complete!) Fsync Behavior (Sync Mode Dependent) Periodic Mode (Default) • Write to commit log @ T=1ms • Send SUCCESS @ T=3ms (before fsync!) • fsync happens @ T=10s (background) Fast! 3ms latency ✓ Risk: 10s window for data loss (but RF=3 protects!) Batch Mode (Safe) • Write to commit log @ T=1ms • fsync immediately @ T=2ms • Send SUCCESS @ T=3ms (after fsync!) Safe! 100% durability ✓ Cost: Slower throughput (fsync per batch = bottleneck)

⏱️ Why Sequential Writes Are So Fast

The physics of sequential I/O:

Random Write (traditional databases):
1. Seek to location (5-10ms - disk head moves!)
2. Wait for rotation (4ms - platter spins)
3. Write data (0.1ms)
Total: ~10-15ms per write

Sequential Write (commit log):
1. Already at end of file (0ms - no seek!)
2. Already at position (0ms - no rotation wait!)
3. Write data (0.1ms)
Total: ~0.1-1ms per write

Result: Sequential is 10-100x faster! This is why commit log can handle 50,000+ writes/sec on a single disk!

⚡ Sync Modes: Periodic vs Batch

Choosing the right sync mode is a critical trade-off between latency and durability!

⚡

Periodic (Default)

Config: commitlog_sync: periodic
Period: 10s default (configurable)
Behavior: fsync every N seconds
Write latency: 3-5ms (no fsync wait!)
Throughput: 50,000+ writes/sec
Risk window: Up to 10s of data at risk
Best for: 99% of production use cases

🛡️

Batch

Config: commitlog_sync: batch
Behavior: fsync BEFORE SUCCESS
Write latency: 8-15ms (includes fsync!)
Throughput: 10,000-20,000 writes/sec
Risk window: Zero! (100% durable)
Cost: 3x slower writes
Best for: Financial, critical data only

Sync Mode Comparison Periodic Mode: Fast, Slight Risk Write 1 3ms Write 2 3ms Write 3 3ms fsync @ 10s mark All 3 writes return in 3ms. fsync happens later (10s). Fast! ✓ Batch Mode: Slower, 100% Safe Write 1 + fsync 12ms Write 2 + fsync 12ms Write 3 + fsync 12ms Each write waits for fsync before SUCCESS. Slower but 100% durable! ✓ Throughput: Periodic = 3 writes in 3ms, Batch = 3 writes in 36ms (12x difference!)
✅

When to Use Periodic

Use case: 99% of production workloads
Why: RF=3 provides redundancy
Risk: If 3 nodes crash in 10s = data loss
Probability: Extremely low (~0.0001%)
Benefit: 3x higher throughput
Examples: User profiles, messages, events

🛡️

When to Use Batch

Use case: Financial, critical transactions
Why: 100% durability guarantee
Risk: Zero data loss window
Cost: 3x slower writes
Worth it when: Cannot afford ANY data loss
Examples: Payments, balances, audit logs

⚖️ The Trade-off Decision

Periodic vs Batch is a classic trade-off:

Periodic (default): Fast (3ms), high throughput (50k writes/sec), tiny risk window (10s), protected by RF=3
Batch: Slower (12ms), lower throughput (15k writes/sec), zero risk, 100% durable

Most production systems use periodic because RF=3 means even if one node crashes mid-fsync, the other 2 replicas have the data safely logged. The probability of all 3 replicas crashing within the same 10s window is astronomically low.

Use batch only for financial transactions, audit logs, or regulatory compliance scenarios where even a 0.0001% risk is unacceptable!

🔄 Crash Recovery: Replaying the Log

When a node crashes, the commit log becomes the source of truth for rebuilding lost memtable data!

1️⃣

Step 1: Detect Replay Needed

On startup: Cassandra scans commit log directory
Check segments: Any with uncommitted writes?
Compare: Segment position vs SSTable markers
Uncommitted data: Writes not yet flushed to SSTable
Decision: If uncommitted data → replay needed
Log message: "Replaying commit log..."

2️⃣

Step 2: Read Mutations

Sequential scan: Read commit log from oldest
Parse binary: Deserialize each mutation
Verify CRC: Checksum validation
Corrupt record: Skip + log warning
Speed: ~1GB/minute (sequential read = fast!)
Progress: Log replay progress percentage

3️⃣

Step 3: Rebuild Memtables

Apply mutations: Re-insert into memtable
Preserve timestamps: Original mutation times kept
Sorted order: Maintain partition + clustering order
All tables: Replay for each affected table
Memory allocation: Memtables recreated
Result: Memtables match pre-crash state

4️⃣

Step 4: Mark Replayed

Track position: Record replay progress
Segment markers: Which segments processed
Completion: All mutations replayed
Log message: "Commit log replay complete"
Metadata update: Mark segments as replayed
Safe to delete: Old segments can be truncated

5️⃣

Step 5: Resume Operations

Node ready: Accept new writes
Memtables active: Can receive mutations
Reads work: Query recovered memtables
Gossip: Node rejoins cluster
Total time: Typically 30-90 seconds
Zero data loss: All ACKed writes recovered ✓

⚠️

Error Handling

CRC mismatch: Skip corrupt record, log warning
Partial write: Discard incomplete mutation
Unknown table: Skip (table may be dropped)
Disk errors: Fail fast, alert operator
Best effort: Recover what's valid
Monitoring: Check logs for corruption warnings

🔄 Netflix Recovery: 20 Nodes, 90 Seconds

Netflix datacenter incident - Cascading node failures:

The Problem:
• Network partition isolates 20 nodes
• Nodes killed after timeout (30 seconds)
• Each node had 1.5GB uncommitted writes
• Total: 30GB of viewing data at risk!

The Recovery:
• Network restored, nodes begin restart
• Each node detects commit log replay needed
• Parallel replay on all 20 nodes
• 1.5GB commit log per node (10 minutes of writes)
• Sequential read: ~1GB/minute
• Replay time: 78 seconds average per node
• Zero corruption detected (CRC checksums passed)
• All 30GB of data recovered perfectly ✓

The Outcome:
• Cluster fully operational in 90 seconds
• Not a single viewing event lost
• Users never noticed the outage
• RF=3 meant other replicas served reads during recovery

Netflix Engineering: "The commit log replay was flawless. 78 seconds to recover 1.5GB of data per node is exactly what we expected from sequential I/O. This is why we trust Cassandra with billions of dollars of viewing data - the durability guarantees actually work in practice!"

⏱️ Recovery Time Estimation

Recovery time formula:

Time = (Commit Log Size) / (Sequential Read Speed)

Typical speeds:
• HDD: ~100 MB/s = 1GB in 10 seconds
• SATA SSD: ~500 MB/s = 1GB in 2 seconds
• NVMe SSD: ~2000 MB/s = 1GB in 0.5 seconds

Example scenarios:
• 100MB commit log on HDD: ~1 second
• 1GB commit log on HDD: ~10 seconds
• 5GB commit log on SATA SSD: ~10 seconds
• 10GB commit log on NVMe: ~5 seconds

Pro tip: Keep commit log under 2GB for fast recovery!

📂 Segments & Truncation: Lifecycle Management

Commit logs are organized into segments that get truncated after memtable flushes!

📄

Segment Structure

Segment size: 32MB default (configurable)
File name: CommitLog-[version]-[timestamp].log
Current segment: Active (receiving writes)
Old segments: Read-only (waiting for truncation)
New segment: Created when current fills
Multiple segments: Normal during high write load

🔄

Segment Rotation

Trigger: Current segment reaches 32MB
Process: Close current, create new
Atomic: No writes lost during rotation
Frequency: Depends on write rate
High load: New segment every few minutes
Low load: Same segment for hours

✂️

Truncation Logic

Trigger: Memtable flushed to SSTable
Check: Which segments have only flushed data?
Safe to delete: All mutations now in SSTable
Delete old: Remove obsolete segments
Keep recent: Unflushed data segments kept
Result: Disk space freed up

💾

Disk Space Management

Typical size: 1-4GB commit log directory
High write load: Can grow to 10-20GB
Flush delays: Disk space can balloon
Monitoring: Alert if >10GB
Emergency: Force flush if disk filling
Best practice: Dedicated disk for commit log

📊

Tracking Flush State

Metadata file: Tracks which segments flushed
Per-table markers: Flush position for each table
Replay markers: Used during crash recovery
Update on flush: Mark segment as flushed
Safe deletion: Only delete if all tables flushed
Efficiency: Minimize replay on restart

⚙️

Configuration Options

commitlog_segment_size_in_mb: 32 default
commitlog_total_space_in_mb: Max total size
Larger segments: Fewer rotations
Smaller segments: Faster truncation
Trade-off: Segment size vs rotation frequency
Recommendation: 32MB good for most cases

⚠️ Commit Log Not Truncating? Debug Tips

Symptom: Commit log directory growing to 50GB+

Common causes:
1. Memtables not flushing: Check memtable_flush_writers (should be >0)
2. Slow compaction: SSTable flush delayed by compaction backlog
3. Disk full: Can't write SSTable = can't truncate commit log
4. Table with no writes: Old segment kept for inactive table
5. Flush errors: Check logs for flush failures

Solutions:
• Manual flush: `nodetool flush` forces memtable flush
• Check disk space: Ensure data disk has room for SSTable
• Monitor compaction: `nodetool compactionstats`
• Increase flush writers: More parallelism = faster flush
• Emergency: Increase commitlog_total_space (temporary!)

🚀 Performance Characteristics

Understanding commit log performance is critical for optimization!

⚡

Write Latency

Periodic mode: ~1ms commit log write
Batch mode: ~2-5ms (includes fsync)
HDD: 1-2ms (sequential = fast!)
SATA SSD: 0.5-1ms
NVMe SSD: 0.1-0.5ms
Bottleneck: Rarely commit log (usually network!)

📈

Throughput Limits

Periodic HDD: 50,000 writes/sec
Periodic SSD: 100,000+ writes/sec
Batch HDD: 10,000-15,000 writes/sec
Batch SSD: 30,000-50,000 writes/sec
Limited by: Sequential write bandwidth
Scale: Add more nodes for more throughput!

💾

Disk Utilization

Space overhead: 1-4GB typical
High load: Can reach 10-20GB
Write bandwidth: 50-300 MB/s
Read (recovery): 100-500 MB/s
I/O pattern: 99% sequential
Efficiency: Very disk-friendly!

📊 Discord Benchmark: Commit Log Impact

Discord tested commit log on same vs separate disk:

Scenario 1: Commit log on same disk as data
• P50 write latency: 8ms
• P99 write latency: 45ms (terrible!)
• Throughput: 25,000 writes/sec
• Problem: Commit log I/O competes with compaction
• Compaction slows down commit log writes
• Latency spikes during compaction

Scenario 2: Commit log on separate NVMe SSD
• P50 write latency: 3ms (62% faster!)
• P99 write latency: 12ms (73% faster!)
• Throughput: 80,000 writes/sec (3.2x higher!)
• No I/O contention with compaction
• Consistent latency (no spikes)
• NVMe sequential write: 500+ MB/s

Result: Discord moved ALL clusters to separate commit log disks. Cost of extra disk: ~$100/node. Performance improvement: 3x throughput, 70% lower latency. ROI: Massive!

🏢 Real Company Commit Log Strategies

💰 Uber: Batch Mode for Financial Accuracy

Use case: Uber Eats payment processing
Requirement: Cannot lose payment data (financial!)

Configuration:
• commitlog_sync: batch
• commitlog_sync_period: N/A (not used in batch)
• Separate commit log disk (RAID 1 mirrored)
• RF=3 across 3 datacenters
• Write CL=QUORUM (2 of 3 DCs)

Results:
• Write latency: 12-15ms (acceptable for payments)
• Throughput: 15,000 payments/sec per node
• Durability: 100% (zero payment loss!)
• Recovery tested monthly (chaos drills)
• Zero data loss in 3 years of production

Trade-off: Slower writes worth it for financial accuracy. 15ms latency acceptable because payments aren't time-critical (user waits anyway). The 100% durability guarantee is mandatory for regulatory compliance!

📱 Apple: Periodic with RF=5

Use case: IoT device telemetry (billions of devices)
Challenge: Extreme write volume (10M+ writes/sec!)

Configuration:
• commitlog_sync: periodic (fast!)
• commitlog_sync_period: 5s (aggressive)
• NVMe SSD for commit log
• RF=5 (high redundancy)
• Write CL=ONE (speed priority)

Strategy:
• Periodic mode for maximum throughput
• 5s fsync period reduces risk window
• RF=5 means 5 independent commit logs
• Probability all 5 nodes crash in 5s: ~0.00001%
• Volume: 10M writes/sec = must optimize!

Results:
• Write latency P50: 2ms
• Throughput: 150,000+ writes/sec per node
• 75,000 nodes = 11 billion writes/sec cluster-wide!
• Data loss events: Zero in 2 years
• RF=5 redundancy sufficient for periodic mode

💬 Discord: Separate Disk + Monitoring

Use case: Message storage (billions of messages)
Optimization: Separate commit log disk

Infrastructure:
• Dedicated NVMe SSD for commit log only
• Separate SATA SSD array for data/SSTables
• No I/O contention between commit log and compaction
• Commit log gets full disk bandwidth

Monitoring:
• Alert if commit log write latency >5ms
• Alert if commit log size >10GB
• Dashboard: Real-time fsync timing
• Chaos testing: Monthly node kills
• Recovery drills: Quarterly

Results:
• P50 latency: 3.2ms (before: 8ms)
• P99 latency: 12ms (before: 45ms)
• Zero latency spikes during compaction
• Throughput: 80,000 msgs/sec per node
• Investment: $100/node for NVMe, 3x performance gain!

✅ Best Practices: Commit Log Optimization

1️⃣

1. Separate Commit Log Disk

Why: Avoid I/O contention with SSTables
Performance gain: 50-70% lower latency
Recommended: NVMe SSD dedicated
Cost: ~$100-300 per node
ROI: Massive (3x throughput typical)
Config: commitlog_directory = /commitlog

2️⃣

2. Use Periodic for Most Cases

Default: commitlog_sync: periodic
Period: 10s default (good balance)
Why safe: RF=3 protects against node crashes
Risk: Minimal with proper RF
Performance: 3x faster than batch
Only use batch: Financial/audit data

3️⃣

3. Monitor Commit Log Latency

Metric: CommitLog write latency
Healthy: <2ms typical
Warning: >5ms investigate
Alert: >10ms disk problem!
Tools: nodetool tablestats, Prometheus
Action: Check disk I/O saturation

4️⃣

4. Test Recovery Regularly

Frequency: Quarterly or monthly
Process: Kill node during write load
Verify: Restart + commit log replay
Check: Zero data loss after recovery
Chaos: Include in chaos engineering
Confidence: Prove durability works!

5️⃣

5. Watch Commit Log Size

Normal: 1-4GB typical
Alert: >10GB indicates problem
Cause: Memtable not flushing
Fix: nodetool flush (manual)
Long-term: Investigate flush delays
Prevention: Monitor flush queue depth

6️⃣

6. Configure Segment Size Wisely

Default: 32MB (good for most)
Larger (64-128MB): Fewer rotations
Smaller (16MB): Faster truncation
Trade-off: Rotation frequency vs recovery time
High volume: 64MB reasonable
Don't: Go below 16MB (too many files!)

🚨 Common Mistakes to Avoid

  • ❌ Commit log on same disk as data: I/O contention kills performance
  • ❌ Using batch mode unnecessarily: 3x slower for minimal benefit (RF=3 already safe!)
  • ❌ Not monitoring commit log size: Can fill disk if memtable flush stuck
  • ❌ Ignoring latency spikes: Early warning of disk degradation
  • ❌ Never testing recovery: Don't wait for production incident!
  • ❌ Disabling commit log: Never do this! Zero durability!

💼 Interview Questions & Answers

1
Explain Cassandra's commit log and why it's critical for durability

Complete Answer:

The commit log is Cassandra's Write-Ahead Log (WAL) - an append-only log that guarantees durability by recording every write to disk before acknowledging success to the client.

How It Works:

  1. Write arrives: Client sends INSERT/UPDATE/DELETE
  2. FIRST step: Write to commit log on disk (sequential append)
  3. THEN: Write to memtable in memory (parallel)
  4. Send SUCCESS: Only after commit log write
  5. Result: If client gets SUCCESS, data is durable on disk

Why It's Critical:

  • Durability guarantee: Memtable is in RAM (volatile), commit log is on disk (durable)
  • Crash recovery: If node crashes before memtable flush, replay commit log on restart
  • Zero data loss: Every acknowledged write survives crashes
  • WAL pattern: Standard database technique (PostgreSQL, MySQL use same approach)

Sequential Append-Only Design:

The commit log uses sequential I/O (append to end of file) which is 100x faster than random I/O:

  • Sequential write: ~1ms on HDD, ~0.5ms on SSD
  • Random write: ~10ms on HDD (disk seek required)
  • Throughput: 100-300 MB/s sequential vs 1-2 MB/s random
  • This is why commit log doesn't kill performance despite being on disk!

Crash Recovery Process:

  1. Node restarts after crash
  2. Cassandra reads commit log segments
  3. Replays all mutations to rebuild memtable
  4. Validates CRC32 checksums (skip corrupt records)
  5. Recovery time: ~1 minute per 1GB of commit log
  6. Result: Memtable reconstructed, zero data loss!

Two Sync Modes:

Periodic (default): fsync every 10s, faster (3ms writes), tiny risk window

Batch: fsync before SUCCESS, slower (12ms writes), 100% durable

Most production uses periodic because RF=3 means other replicas have the data even if one node crashes mid-fsync.

Why Without Commit Log = Disaster:

If Cassandra had no commit log:

  • Every crash = lose all memtable data (gigabytes!)
  • Users would see writes disappear randomly
  • No durability guarantee possible
  • Database would be unreliable for production

Key Takeaway: The commit log is the foundation of Cassandra's durability. It pays a small performance cost (~1ms per write) but guarantees that acknowledged writes survive ANY crash. This is what makes Cassandra production-grade!

2
Compare periodic vs batch sync modes - when to use each?

Complete Answer:

Periodic and batch are two sync modes that trade off between write latency and durability guarantees.

Periodic Mode (Default):

  • Behavior: Write to commit log, send SUCCESS immediately, fsync every N seconds (10s default)
  • Write latency: 3-5ms (no fsync wait)
  • Throughput: 50,000+ writes/sec per node
  • Risk window: Up to 10s of data at risk if node crashes before fsync
  • Probability of loss: Very low with RF=3 (all 3 nodes would need to crash in same 10s window)

Batch Mode:

  • Behavior: Write to commit log, fsync immediately, THEN send SUCCESS
  • Write latency: 8-15ms (includes fsync overhead)
  • Throughput: 10,000-20,000 writes/sec per node
  • Risk window: ZERO - 100% durable before client sees SUCCESS
  • Cost: 3-5x slower writes due to fsync overhead

When to Use Periodic (99% of cases):

  • User-facing data (profiles, messages, posts)
  • High-volume workloads (need throughput)
  • When RF>=3 (redundancy protects against single node crash)
  • When 10s risk window acceptable
  • Example: Discord messages, Netflix viewing history, social media posts

When to Use Batch (rare cases):

  • Financial transactions (payments, balances)
  • Audit logs (regulatory requirement)
  • When even 0.001% data loss unacceptable
  • When write volume low enough to handle 3x slower performance
  • Example: Uber Eats payments, banking transactions

Performance Comparison:

Benchmark: 100k writes, RF=3, single DC

Periodic: Total time 2 seconds, P50 latency 3ms, P99 latency 8ms

Batch: Total time 7 seconds, P50 latency 12ms, P99 latency 25ms

Batch is 3.5x slower!

Why Periodic is Safe with RF=3:

For data loss with periodic mode:

  1. Write happens at T=0, logged to commit log (not yet fsynced)
  2. Node crashes at T=5s (before 10s fsync)
  3. BUT: Write was also sent to 2 other replicas (RF=3)
  4. Those 2 replicas also have the write in their commit logs
  5. For data loss: All 3 nodes must crash in same 10s window
  6. Probability: ~0.0001% (extremely rare!)

Key Insight: Periodic mode isn't "unsafe" - it's a calculated risk that's protected by replication. Batch mode is only needed when that tiny risk is unacceptable (financial data). The 3x performance cost of batch mode is rarely worth it!

Advertisement

Responsive Ad