Node Failure Handling
Master Cassandra's fault tolerance with animated diagrams showing failure detection, hinted handoff, and automatic recovery!
📖 Instagram's Hardware Failure: Zero Downtime
In 2015, Instagram's Cassandra cluster experienced a catastrophic hardware failure - 12 nodes went offline simultaneously. Result? ZERO DOWNTIME! 400+ million users noticed nothing. Cassandra's automatic failure detection, hinted handoff, and read repair kept everything running seamlessly!
🔍 Failure Detection - Phi Accrual Algorithm
How Phi (Φ) Works
- Φ = 0: Healthy (regular heartbeats)
- Φ = 3-5: Suspicious (slight delays)
- Φ > 8: Failed (stopped responding)
- Adaptive: Automatically adjusts to network conditions
- Detection Time: 2-3 seconds typical
Real-World Detection Example: Netflix
Scenario: Netflix's US-East datacenter node crashes during peak streaming hours (8 PM Friday).
Timeline:
- T=0.0s: Node 10.1.0.5 suffers hardware failure, stops responding
- T=1.0s: Neighboring nodes miss first heartbeat, Φ = 3 (suspicious)
- T=2.0s: Second heartbeat missed, Φ = 6 (highly suspicious)
- T=3.2s: Φ crosses threshold (8), node marked DOWN
- T=3.5s: Gossip spreads failure state to 50% of cluster
- T=5.0s: Entire 300-node cluster aware of failure
- T=5.1s: Client drivers reroute traffic, hinted handoff activated
✅ Result: 80 million concurrent viewers experience zero disruption!
Live Console: Monitoring Failure Detection
# Check node status in real-time
$ nodetool status
Datacenter: US-East
===================
Status=Up/Down
|/ State=Normal/Leaving/Joining/Moving
-- Address Load Tokens Owns Host ID Rack
UN 10.1.0.1 245.5 GB 256 33.3% a1b2c3d4-e5f6-7890-abcd-ef1234567890 rack1
UN 10.1.0.2 238.2 GB 256 33.2% b2c3d4e5-f6g7-8901-bcde-f12345678901 rack1
UN 10.1.0.3 251.8 GB 256 33.4% c3d4e5f6-g7h8-9012-cdef-012345678912 rack2
DN 10.1.0.4 247.3 GB 256 33.3% d4e5f6g7-h8i9-0123-defg-123456789123 rack2
↑
DOWN - Failure detected!
# Check gossip info for failed node
$ nodetool gossipinfo | grep 10.1.0.4
/10.1.0.4
generation:1702941234
heartbeat:999
STATUS:DOWN
LOAD:247.3GB
SCHEMA:59adb24e-1234-5678-90ab-cdef12345678
DC:US-East
RACK:rack2
RELEASE_VERSION:4.1.3
PHI:12.5 ← Suspicion level above threshold!
# Watch live gossip updates (updates every second)
$ watch -n 1 'nodetool gossipinfo | grep -A 3 10.1.0.4'
# Check when node was marked down
$ grep "10.1.0.4.*DOWN" /var/log/cassandra/system.log
INFO [GossipStage:1] 2024-12-26 20:15:33,215 Gossiper.java:1234
- InetAddress /10.1.0.4 is now DOWN
INFO [GossipStage:1] 2024-12-26 20:15:33,216 StorageService.java:5678
- Node /10.1.0.4 state changed to DOWN, phi = 12.5
Configuring Phi Threshold
Tuning Phi Accrual Settings
Default threshold is phi_convict_threshold = 8 in cassandra.yaml
# cassandra.yaml - Adjust failure detection sensitivity # Default (recommended for most environments) phi_convict_threshold: 8 # More sensitive (faster detection, more false positives) phi_convict_threshold: 5 # Use case: Low-latency networks, quick failover needed # Less sensitive (slower detection, fewer false positives) phi_convict_threshold: 12 # Use case: High-latency/unstable networks, avoid flapping # Check current setting $ grep phi_convict_threshold /etc/cassandra/cassandra.yaml phi_convict_threshold: 8
⚠️ Warning: Lower threshold = faster detection but more false positives during GC pauses or network hiccups!
⚡ What Happens When Node Fails
Detailed Breakdown: Request Handling During Failure
Write Request Processing (RF=3, CL=QUORUM)
Example: Client writes user profile update, one replica is down
# Client executes write
INSERT INTO users (user_id, name, email)
VALUES (12345, 'Alice', 'alice@example.com')
USING CONSISTENCY QUORUM;
# Behind the scenes:
1. Coordinator (Node 10.1.0.1) receives write
2. Determines replicas: 10.1.0.2, 10.1.0.3, 10.1.0.4
3. Sends write to all 3 replicas in parallel
→ 10.1.0.2: SUCCESS (50ms) ✅
→ 10.1.0.3: SUCCESS (45ms) ✅
→ 10.1.0.4: TIMEOUT (node is DOWN) ❌
4. QUORUM = 2, got 2 responses → Write SUCCEEDS
5. Coordinator stores hint for 10.1.0.4:
{
target_node: 10.1.0.4,
data: ,
timestamp: 1703012345678,
ttl: 10800000 // 3 hours in ms
}
6. Return success to client (total latency: 52ms)
Read Request Processing (RF=3, CL=QUORUM)
Example: Client reads user profile, one replica is down
# Client executes read
SELECT name, email FROM users
WHERE user_id = 12345
USING CONSISTENCY QUORUM;
# Behind the scenes:
1. Coordinator (Node 10.1.0.1) receives read
2. Determines replicas: 10.1.0.2, 10.1.0.3, 10.1.0.4
3. Sends read to all 3 replicas (for read repair comparison)
→ 10.1.0.2: {name: 'Alice', email: 'alice@example.com'} (15ms) ✅
→ 10.1.0.3: {name: 'Alice', email: 'alice@example.com'} (18ms) ✅
→ 10.1.0.4: TIMEOUT (node is DOWN) ❌
4. QUORUM = 2, got 2 responses → Read SUCCEEDS
5. Compare timestamps (both identical) → No repair needed
6. Return data to client (total latency: 19ms)
# If data was inconsistent:
→ 10.1.0.2: {name: 'Alice', ts: 1703012345678} ✅ LATEST
→ 10.1.0.3: {name: 'Bob', ts: 1703012340000} ⚠️ STALE
Coordinator:
- Returns 'Alice' to client (latest)
- Background: Sends 'Alice' update to 10.1.0.3
- Read repair completes silently
Live Console: Monitoring Active Failures
# Check current failures and unavailable exceptions
$ nodetool tpstats | grep -i unavailable
Pool Name Active Pending Completed
MutationStage 2 15 95847562
ReadStage 0 0 48293847
RequestResponseStage 0 0 12847293
# Check if cluster is experiencing timeouts
$ grep "Timeout" /var/log/cassandra/system.log | tail -20
WARN [ReadStage-2] 2024-12-26 20:15:35 - Read timeout for /10.1.0.4
INFO [RequestResponseStage-3] 2024-12-26 20:15:35 -
Using replica /10.1.0.2 instead of /10.1.0.4
# Monitor hints being stored
$ nodetool getendpoints mykeyspace users 12345
10.1.0.2
10.1.0.3
10.1.0.4 ← This replica is down, hints being stored
# Check hint statistics
$ nodetool handoff stats
Total hints: 15847
Hints per host:
10.1.0.4: 15847 hints (DOWN node)
Oldest hint: 124 seconds ago
Newest hint: 2 seconds ago
# Real-time hint monitoring
$ watch -n 5 'nodetool handoff stats'
# Check coordinator load
$ nodetool status | grep 10.1.0.1
UN 10.1.0.1 248.7 GB 256 33.3% ... rack1
↑ Load increasing due to hints storage
Real Scenario: Uber's Multi-Node Failure
Incident: Uber's ride-matching Cassandra cluster lost 3 nodes simultaneously due to network switch failure
Cluster Configuration:
- 12-node cluster, RF=3, CL=QUORUM
- Nodes: 10.1.0.1 through 10.1.0.12
- Failed nodes: 10.1.0.4, 10.1.0.5, 10.1.0.6 (all on same switch)
Impact Analysis:
- Token Distribution: Each node owns ~8.3% of data
- Affected Keys: ~25% of keys had one replica down
- Still Available: Every key had 2 replicas up (QUORUM met!)
- Hint Load: 9 coordinator nodes storing hints
- Write Latency: Increased 5ms average (waiting for timeouts)
- Read Latency: Unchanged (healthy replicas responded quickly)
✅ Result: 100% of ride requests processed successfully! Users noticed 0.005s latency increase only!
# Monitoring during incident $ nodetool status DN 10.1.0.4 247GB 256 8.3% (3 minutes ago) DN 10.1.0.5 245GB 256 8.3% (3 minutes ago) DN 10.1.0.6 251GB 256 8.3% (3 minutes ago) $ nodetool handoff stats Total hints: 284,592 10.1.0.4: 94,864 hints 10.1.0.5: 95,128 hints 10.1.0.6: 94,600 hints Hint storage: 1.2 GB total across coordinators Growth rate: ~800 hints/second # After network switch restored (45 minutes later) $ nodetool status UN 10.1.0.4 247GB 256 8.3% ← Back UP! UN 10.1.0.5 245GB 256 8.3% ← Back UP! UN 10.1.0.6 251GB 256 8.3% ← Back UP! # Hints automatically replaying $ grep "Replaying" /var/log/cassandra/system.log INFO - Replaying 94,864 hints for /10.1.0.4 INFO - Replaying 95,128 hints for /10.1.0.5 INFO - Replaying 94,600 hints for /10.1.0.6 # 8 minutes later - all hints replayed $ nodetool handoff stats Total hints: 0 All nodes synchronized!
📬 Hinted Handoff - Zero Data Loss
Hint Limitations
- TTL: Hints expire after 3 hours (default)
- Best Effort: Not guaranteed (coordinator could fail)
- Long Outages: Need to run nodetool repair
- Storage: Stored on coordinator's local disk
Detailed Hint Storage Mechanics
Hint Storage Structure
Hints are stored in the system.hints table on the coordinator node
# Examine hints table structure
$ cqlsh -e "DESC TABLE system.hints;"
CREATE TABLE system.hints (
target_id uuid, -- UUID of the down node
hint_id timeuuid, -- Unique hint identifier
message_version int, -- Protocol version
mutation blob, -- The actual write data
PRIMARY KEY (target_id, hint_id, message_version)
);
# Check hint files on disk
$ ls -lh /var/lib/cassandra/hints/
-rw-r--r-- 1 cassandra 2.4M hints-d4e5f6g7-h8i9.log
-rw-r--r-- 1 cassandra 1.8M hints-d4e5f6g7-h8i9-1.log
↑
Target node ID
# View hint statistics
$ nodetool getendpointstates | grep HINTS
HINTS[d4e5f6g7-h8i9-0123-defg-123456789123]: 15847
# Configuration in cassandra.yaml
max_hint_window_in_ms: 10800000 # 3 hours
max_hints_delivery_threads: 2 # Parallel replay threads
hints_directory: /var/lib/cassandra/hints
hints_flush_period_in_ms: 10000 # Flush to disk every 10s
max_hints_file_size_in_mb: 128 # Max file size
Live Console: Monitoring Hint Replay
# Real-time hint monitoring during node recovery $ watch -n 2 'nodetool handoff stats' # Before node returns Total hints: 15847 Hints per host: 10.1.0.4: 15847 hints Oldest hint: 1847 seconds ago (30 minutes) Hint storage: 245 MB # Node 10.1.0.4 comes back online at T=0 # Coordinator detects and begins replay immediately # T=10 seconds - Replay starting Total hints: 15847 Currently replaying to: 10.1.0.4 Replay rate: 1250 hints/second Estimated completion: 12 seconds # T=15 seconds - Active replay Total hints: 9823 Currently replaying to: 10.1.0.4 Replay rate: 1205 hints/second Estimated completion: 8 seconds # T=23 seconds - Completed! Total hints: 0 All hints successfully replayed to 10.1.0.4 # Check logs for replay confirmation $ tail -f /var/log/cassandra/system.log | grep -i hint INFO [HintsDispatcher:1] HintService.java:345 - Started hint dispatch to /10.1.0.4 INFO [HintsDispatcher:1] HintService.java:456 - Dispatched 15847 hints to /10.1.0.4 in 12.3 seconds INFO [HintsDispatcher:1] HintService.java:467 - Deleted hint file hints-d4e5f6g7-h8i9.log after successful replay # Verify node caught up $ nodetool describecluster | grep "Schema versions" Schema versions: 59adb24e-1234-5678-90ab-cdef12345678: [10.1.0.1, 10.1.0.2, 10.1.0.3, 10.1.0.4] ↑ All nodes have same schema - synchronized!
Real Scenario: Twitter's Hint Overflow
Problem: Twitter node went down during major event (World Cup final), 4.5 hours outage
What Happened:
- T=0: Node 10.1.0.7 crashes during peak load
- T=0-3 hours: Coordinators store hints (working perfectly)
- T=3 hours: Hints reach TTL, start expiring
- T=3-4.5 hours: New writes NOT hinted (window expired)
- T=4.5 hours: Node returns, hints from first 3 hours replay
- Result: Data from hours 3-4.5 MISSING from this replica
# T=3 hours - Hints expiring $ nodetool handoff stats Total hints: 284,592 Expiring hints detected! 10.1.0.7: 284,592 hints Oldest: 10802 seconds (3 hours 2 seconds) - EXPIRED WARN: Hints for /10.1.0.7 exceeding max_hint_window_in_ms $ grep "hints.*expir" /var/log/cassandra/system.log WARN [HintedHandoff:1] - Hints for /10.1.0.7 are expiring WARN [HintedHandoff:1] - Deleted 45,284 expired hints for /10.1.0.7 # T=4.5 hours - Node returns $ nodetool status UN 10.1.0.7 ← Back up! # Check hint replay $ nodetool handoff stats Replaying remaining hints to 10.1.0.7 Replayed: 239,308 hints (hints from first 3 hours) Missing: ~45,284 writes (from hours 3-4.5) # Solution: Run repair to catch up missing data $ nodetool repair -pr mykeyspace [2024-12-26 20:15:45] Starting repair for keyspace mykeyspace [2024-12-26 20:18:23] Repair session completed [2024-12-26 20:18:23] Streamed 2.3 GB to /10.1.0.7 All replicas now in sync!
⚠️ Lesson: For outages > 3 hours, always run nodetool repair!
Configuring Hint Behavior
# cassandra.yaml - Hint configuration options # Maximum time to store hints (default 3 hours) max_hint_window_in_ms: 10800000 # Increase for longer expected outages: max_hint_window_in_ms: 21600000 # 6 hours # Maximum hints in memory before flushing to disk max_hints_in_progress: 128 # Number of threads for hint replay max_hints_delivery_threads: 2 # Increase for faster replay: max_hints_delivery_threads: 4 # Throttle hint replay to avoid overwhelming recovered node hinted_handoff_throttle_in_kb: 1024 # 1 MB/sec per delivery thread # Faster replay: hinted_handoff_throttle_in_kb: 10240 # 10 MB/sec # Disable hints entirely (NOT recommended for production!) # hinted_handoff_enabled: false # Check current settings $ nodetool getcompactionthreshold hints Current compaction threshold: min=4, max=32 # Monitor hint replay performance $ nodetool proxyhistograms | grep -A 5 "Hint" Hint write latencies (microseconds): 50%: 125 95%: 500 99%: 2000
🔧 Read Repair - Automatic Consistency
Read Repair Types
- Foreground: During reads with CL > ONE (automatic)
- Background: Periodic random reads (configurable %)
- Anti-Entropy: Manual repair via nodetool
Configuring Read Repair
# Configure read repair chance per table
CREATE TABLE users (
user_id int PRIMARY KEY,
name text,
email text
) WITH read_repair_chance = 0.1 -- 10% of reads trigger full comparison
AND dclocal_read_repair_chance = 0.1; -- Local DC only
# Check current read repair settings
$ cqlsh -e "SELECT table_name, read_repair_chance, dclocal_read_repair_chance
FROM system_schema.tables
WHERE keyspace_name = 'mykeyspace';"
table_name | read_repair_chance | dclocal_read_repair_chance
-----------+--------------------+---------------------------
users | 0.1 | 0.1
sessions | 0.0 | 0.1
# Modify read repair chance
ALTER TABLE users WITH read_repair_chance = 0.2;
# Disable read repair (use only anti-entropy repair)
ALTER TABLE users WITH read_repair_chance = 0.0
AND dclocal_read_repair_chance = 0.0;
Live Console: Monitoring Read Repair
# Check read repair statistics $ nodetool tablehistograms mykeyspace users | grep -A 10 "Read Repair" Read Repair Attempts: 15,847 Read Repair Repaired: 2,341 (14.8% found inconsistencies) Read Repair Background: 1,584 # View read repair in logs $ grep "ReadRepair" /var/log/cassandra/debug.log | tail -20 DEBUG [ReadRepairStage:42] ReadRepairHandler.java:234 - Read repair for key user123: replicas [10.1.0.2, 10.1.0.3, 10.1.0.4] DEBUG [ReadRepairStage:42] ReadRepairHandler.java:267 - Digest mismatch detected for user123 10.1.0.2: digest=a1b2c3d4 (timestamp: 1703012345678) 10.1.0.3: digest=e5f6g7h8 (timestamp: 1703012340000) ← STALE 10.1.0.4: digest=a1b2c3d4 (timestamp: 1703012345678) DEBUG [ReadRepairStage:42] ReadRepairHandler.java:289 - Sending repair mutation to 10.1.0.3 DEBUG [ReadRepairStage:42] ReadRepairHandler.java:312 - Read repair completed for user123 in 45ms # Monitor read repair rate $ watch -n 5 'nodetool tpstats | grep ReadRepair' Pool Name Active Pending Completed ReadRepairStage 2 15 847293 AntiEntropyStage 0 0 12847 # Check if read repair is working effectively $ nodetool tablestats mykeyspace.users | grep -i bloom Bloom filter false positives: 0 Bloom filter false ratio: 0.00000 Bloom filter space used: 14523 KB Read repair: 2341 repairs, 14.8% of reads ↑ This tells you how often inconsistencies are found
Real Scenario: Spotify's Eventual Consistency
Scenario: Spotify user updates playlist, read repair catches stale replica
# T=0: User adds song to playlist UPDATE playlists SET songs = songs + ['song_id_789'] WHERE user_id = 12345 AND playlist_id = 'workout' USING CONSISTENCY QUORUM; Write goes to 3 replicas: 10.1.0.2: SUCCESS (ts: 1703012345678) ✅ 10.1.0.3: SUCCESS (ts: 1703012345678) ✅ 10.1.0.4: TIMEOUT → Hint stored # T=30s: User reads playlist from different datacenter SELECT songs FROM playlists WHERE user_id = 12345 AND playlist_id = 'workout' USING CONSISTENCY QUORUM; Coordinator queries all replicas: 10.1.0.2: [song1, song2, song_id_789] (ts: 1703012345678) 10.1.0.3: [song1, song2, song_id_789] (ts: 1703012345678) 10.1.0.4: [song1, song2] (ts: 1703012340000) ← STALE! Read Repair Process: 1. Coordinator compares digests 2. Detects 10.1.0.4 is stale (older timestamp) 3. Returns latest to user immediately: [song1, song2, song_id_789] 4. Background: Sends full row to 10.1.0.4 5. 10.1.0.4 updates: [song1, song2, song_id_789] (ts: 1703012345678) Result: All replicas synchronized within 50ms!
Anti-Entropy Repair (nodetool repair)
When to Use Anti-Entropy Repair
Full repair using Merkle trees - comprehensive but resource-intensive
- Weekly Schedule: Best practice for production
- After Long Outages: Node down > 3 hours
- After Deletes: Before gc_grace_seconds expires
- Data Migrations: Adding/removing nodes
- Suspected Issues: Seeing stale data
# Basic repair - repairs primary ranges only (RECOMMENDED)
$ nodetool repair -pr mykeyspace
[2024-12-26 20:15:45] Starting repair command
[2024-12-26 20:15:46] Repair session id: 5e8a7c90-a432-11ef-b123-0242ac130003
[2024-12-26 20:15:47] Building Merkle trees for mykeyspace.users
[2024-12-26 20:15:58] Merkle tree build completed in 11.2s
[2024-12-26 20:15:58] Comparing Merkle trees with replicas
[2024-12-26 20:16:03] Found 847 inconsistent ranges
[2024-12-26 20:16:04] Streaming 245 MB to /10.1.0.4
[2024-12-26 20:18:45] Repair completed for mykeyspace
# Full repair (all replicas, heavy on resources)
$ nodetool repair mykeyspace users
# Incremental repair (only changed data since last repair)
$ nodetool repair -inc mykeyspace
# Repair specific token range
$ nodetool repair -pr -st -9223372036854775808 -et -9223372036854775000 mykeyspace
# Monitor active repairs
$ nodetool compactionstats
pending tasks: 0
- 10.1.0.4:
Active compaction remaining time : n/a
- 10.1.0.2:
Active repair remaining time : 42m15s
# Check repair history
$ nodetool repair_admin list --all
Repair sessions:
id: 5e8a7c90 state: COMPLETED ranges: 847 progress: 100%
id: 6f9b8d01 state: IN_PROGRESS ranges: 456 progress: 67%
# Schedule automatic repairs (cron)
$ crontab -e
# Run repair every Sunday at 2 AM
0 2 * * 0 /usr/bin/nodetool repair -pr mykeyspace >> /var/log/cassandra-repair.log 2>&1
Real-World Repair Example: LinkedIn
Challenge: LinkedIn runs weekly repairs on 500+ node cluster
Strategy:
- Staggered Schedule: 1 node per hour (500 nodes = 21 days cycle)
- Off-Peak Hours: Run during low traffic (2-6 AM local time)
- Incremental Mode: Only repair changed data (-inc flag)
- Monitoring: Track completion, failures, data streamed
- Automation: Cassandra Reaper tool for orchestration
# LinkedIn's automated repair script
#!/bin/bash
# repair-node.sh - Run on each node via cron
KEYSPACE="linkedin_prod"
LOG="/var/log/cassandra/repair-$(date +%Y%m%d).log"
echo "Starting repair at $(date)" >> $LOG
# Check cluster health before starting
if nodetool status | grep -q "DN"; then
echo "ERROR: Cluster has down nodes, skipping repair" >> $LOG
exit 1
fi
# Check disk space
DISK_FREE=$(df -h /var/lib/cassandra | awk 'NR==2 {print $5}' | sed 's/%//')
if [ $DISK_FREE -gt 80 ]; then
echo "ERROR: Disk usage >80%, skipping repair" >> $LOG
exit 1
fi
# Run incremental repair on primary ranges
nodetool repair -pr -inc $KEYSPACE >> $LOG 2>&1
if [ $? -eq 0 ]; then
echo "Repair completed successfully at $(date)" >> $LOG
# Send success metric to monitoring
echo "cassandra.repair.success:1|c" | nc -u -w1 localhost 8125
else
echo "ERROR: Repair failed at $(date)" >> $LOG
# Alert on-call
curl -X POST "https://alerts.linkedin.com/cassandra-repair-failed" \
-d "node=$(hostname)&keyspace=$KEYSPACE"
fi
# Typical output
[2024-12-26 02:00:15] Starting repair
[2024-12-26 02:00:18] Incremental repair using saved session
[2024-12-26 02:00:45] Merkle tree diff: 234 ranges need repair
[2024-12-26 02:01:12] Streaming 45 MB to replicas
[2024-12-26 02:03:47] Repair completed in 3m32s
[2024-12-26 02:03:47] Data streamed: 45 MB
[2024-12-26 02:03:47] Ranges repaired: 234/8192 (2.9%)
🔧 Troubleshooting Node Failures
Real-world troubleshooting scenarios with console commands and solutions.
Scenario 1: Node Won't Rejoin After Restart
# Problem: Node stuck in DN state after restart
$ nodetool status
DN 10.1.0.4 247GB 256 8.3% (10 minutes ago)
# Step 1: Check if Cassandra process is running
$ systemctl status cassandra
● cassandra.service - Cassandra
Active: active (running) since 2024-12-26 20:15:00
$ ps aux | grep cassandra
cassandra 12345 250 15.5 Process is running!
# Step 2: Check logs for errors
$ tail -100 /var/log/cassandra/system.log
ERROR [main] - Fatal configuration error
ERROR [main] - seed_provider parameters list is empty
↑ Problem found: No seeds configured!
# Step 3: Fix cassandra.yaml
$ vim /etc/cassandra/cassandra.yaml
seed_provider:
- class_name: org.apache.cassandra.locator.SimpleSeedProvider
parameters:
- seeds: "10.1.0.1,10.1.0.2" ← Added missing seeds
# Step 4: Restart Cassandra
$ systemctl restart cassandra
# Step 5: Verify node joins
$ watch -n 2 'nodetool status | grep 10.1.0.4'
UN 10.1.0.4 247GB 256 8.3% ← Back UP! ✅
# Step 6: Check gossip
$ nodetool gossipinfo | grep 10.1.0.4 -A 5
/10.1.0.4
STATUS:NORMAL
LOAD:247GB
generation:1703012456
heartbeat:55 ← Heartbeats flowing!
Scenario 2: Hints Piling Up, Disk Full
# Problem: Hints directory consuming too much disk
$ df -h /var/lib/cassandra
Filesystem Size Used Avail Use% Mounted on
/dev/sdb1 500G 475G 25G 95% /var/lib/cassandra
↑ Critical!
$ du -sh /var/lib/cassandra/hints/
45G /var/lib/cassandra/hints/
↑ Hints taking 45GB!
$ nodetool handoff stats
Total hints: 4,847,592
Hints per host:
10.1.0.4: 4,847,592 hints (DOWN for 6 hours!)
Oldest hint: 21,847 seconds ago (6+ hours)
⚠️ Exceeding max_hint_window_in_ms!
# Solution 1: Bring the down node back (best option)
# If node is permanently gone or can't be recovered:
# Solution 2: Remove dead node from cluster
$ nodetool assassinate 10.1.0.4
Node 10.1.0.4 assassinated - gossip state removed
# Solution 3: Delete hints manually (last resort!)
$ nodetool truncatehints 10.1.0.4
Truncated hints for /10.1.0.4
# Verify disk space recovered
$ du -sh /var/lib/cassandra/hints/
128M /var/lib/cassandra/hints/ ← Much better!
$ df -h /var/lib/cassandra
Filesystem Size Used Avail Use% Mounted on
/dev/sdb1 500G 430G 70G 86% /var/lib/cassandra
↑ Back to safe levels
# Solution 4: Run repair on the cluster
# Since hints were deleted, need to resync data
$ for node in $(nodetool status | grep UN | awk '{print $2}'); do
echo "Repairing $node..."
nodetool repair -pr mykeyspace -h $node
done
Scenario 3: Split Brain - Cluster Partition
# Problem: Network partition splits cluster into 2 groups
# Datacenter 1 can't reach Datacenter 2
# From DC1 node:
$ nodetool status
Datacenter: US-East (DC1)
UN 10.1.0.1 245GB 256 16.7%
UN 10.1.0.2 238GB 256 16.6%
UN 10.1.0.3 251GB 256 16.8%
DN 10.2.0.1 0GB 256 16.6% ← DC2 appears down
DN 10.2.0.2 0GB 256 16.6% ← DC2 appears down
DN 10.2.0.3 0GB 256 16.7% ← DC2 appears down
# From DC2 node:
$ nodetool status
Datacenter: EU-West (DC2)
DN 10.1.0.1 0GB 256 16.7% ← DC1 appears down
DN 10.1.0.2 0GB 256 16.6% ← DC1 appears down
DN 10.1.0.3 0GB 256 16.8% ← DC1 appears down
UN 10.2.0.1 247GB 256 16.6%
UN 10.2.0.2 242GB 256 16.6%
UN 10.2.0.3 249GB 256 16.7%
# Check network connectivity
$ ping 10.2.0.1
Request timeout ← Network is down!
$ traceroute 10.2.0.1
1 gateway (10.1.0.254) 1.2 ms
2 * * * ← Packets lost at inter-DC router
# Impact on operations:
# - LOCAL_QUORUM: Works in each DC independently ✅
# - QUORUM: Fails (can't reach majority of 6 nodes) ❌
# - LOCAL_ONE: Works in each DC ✅
# Check write failures
$ tail -f /var/log/cassandra/system.log | grep Unavailable
UnavailableException: Cannot achieve consistency QUORUM
(4 required but only 3 alive)
# Emergency: Switch to LOCAL_QUORUM
# In application code:
session.execute(
statement,
consistency_level=ConsistencyLevel.LOCAL_QUORUM # Use local DC only
);
# After network restored:
$ ping 10.2.0.1
Reply from 10.2.0.1: time=45ms ← Network back!
$ nodetool status
All nodes show UN ← Cluster reunited!
# Run repair to fix diverged data
$ nodetool repair -pr mykeyspace
Repairing data differences from split-brain period
Scenario 4: Slow Queries During Node Failure
# Problem: Application experiencing slow queries after node failure
# Check current performance
$ nodetool tablestats mykeyspace.users | grep -i latency
Local read latency: 125.5 ms (up from 15ms normal!)
↑ 8x increase!
# Investigate
$ nodetool tpstats
Pool Name Active Pending Blocked
MutationStage 15 458 0 ← High pending!
ReadStage 8 234 0 ← Backlog
RequestResponseStage 2 45 0
Dropped Messages:
READ : 847
MUTATION : 234
# Check for timeouts in logs
$ grep "timeout" /var/log/cassandra/system.log | wc -l
4,592 ← Many timeouts!
$ grep "timeout" /var/log/cassandra/system.log | tail -5
WARN [ReadStage-4] - Read timeout: waited 5000ms for /10.1.0.4
WARN [ReadStage-7] - Read timeout: waited 5000ms for /10.1.0.4
↑ Consistently timing out to 10.1.0.4 (dead node)
# Check node status
$ nodetool status | grep 10.1.0.4
DN 10.1.0.4 247GB 256 8.3% ← Node is down
# Problem: Client driver still trying to contact dead node!
# Solution 1: Update driver connection settings (application side)
# Python driver example:
cluster = Cluster(
['10.1.0.1', '10.1.0.2', '10.1.0.3'],
load_balancing_policy=TokenAwarePolicy(RoundRobinPolicy()),
reconnection_policy=ConstantReconnectionPolicy(delay=5.0)
)
# Driver will automatically:
# - Mark 10.1.0.4 as down after connection failures
# - Stop routing requests to it
# - Try reconnecting every 5 seconds
# Solution 2: Check driver is token-aware
# Verify in application logs:
INFO - Cluster topology: {
'10.1.0.1': UP (latency: 15ms),
'10.1.0.2': UP (latency: 18ms),
'10.1.0.3': UP (latency: 16ms),
'10.1.0.4': DOWN (marked 300 seconds ago)
}
# Solution 3: Monitor recovery
$ watch -n 5 'nodetool tablestats mykeyspace.users | grep "Local read latency"'
T=0: Local read latency: 125.5 ms (with timeouts)
T=10s: Local read latency: 45.2 ms (driver adapted)
T=30s: Local read latency: 18.3 ms (back to normal!) ✅
# Solution 4: If node is permanently gone, remove it
$ nodetool removenode d4e5f6g7-h8i9-0123-defg-123456789123
Removing node, streaming ranges to remaining replicas...
Node removed successfully!
$ nodetool cleanup
Running cleanup on all nodes to remove old token ranges
Scenario 5: Lost Quorum - Emergency Recovery
# Disaster: 2 of 3 nodes down, lost QUORUM!
$ nodetool status
Datacenter: Production
UN 10.1.0.1 247GB 256 33.3% ← Only survivor
DN 10.1.0.2 238GB 256 33.2% ← Down
DN 10.1.0.3 251GB 256 33.4% ← Down
QUORUM = 2, but only 1 node up → OPERATIONS FAILING!
# All queries fail:
$ cqlsh -e "SELECT * FROM users WHERE user_id=123" \
--request-timeout=10
UnavailableException: Cannot achieve consistency level QUORUM
(2 required but only 1 alive)
# Emergency Option 1: Lower consistency (TEMPORARY!)
$ cqlsh
cqlsh> CONSISTENCY ONE;
Consistency level set to ONE.
cqlsh> SELECT * FROM users WHERE user_id=123;
user_id | name | email
---------+-------+----------------
123 | Alice | alice@test.com
✅ Works with CL=ONE
⚠️ WARNING: Data may be stale! Only 1 replica responding!
# Emergency Option 2: Quickly restore at least 1 more node
# Check why nodes are down:
$ ssh 10.1.0.2
$ systemctl status cassandra
● cassandra.service - Cassandra
Active: failed (Result: exit-code)
$ journalctl -u cassandra | tail -50
OutOfMemoryError: Java heap space
↑ JVM crashed due to memory!
# Quick fix: Increase heap and restart
$ vim /etc/cassandra/jvm.options
-Xms16G # Increased from 8G
-Xmx16G # Increased from 8G
$ systemctl start cassandra
# Monitor recovery
$ watch -n 2 'nodetool status'
T=0: DN 10.1.0.2
T=30s: UJ 10.1.0.2 ← Joining
T=60s: UN 10.1.0.2 ← UP! QUORUM restored! ✅
# Now 2 of 3 nodes up = QUORUM works again!
$ cqlsh
cqlsh> CONSISTENCY QUORUM;
Consistency level set to QUORUM.
cqlsh> SELECT * FROM users WHERE user_id=123;
✅ Works with QUORUM again!
# After Crisis: Root cause analysis
# Why did 2 nodes fail together?
# - Same heap configuration (8GB too small)
# - Large compaction triggered OOM on both
# - Need better capacity planning
# Preventive measures:
1. Increase heap size appropriately (16GB)
2. Configure max_heap_size in cassandra.yaml
3. Set up memory monitoring alerts
4. Implement circuit breakers in application
5. Add more nodes to increase fault tolerance (RF=5)
✅ Best Practices for Failure Handling
Set RF=3 Minimum
Replication Factor 3 allows tolerating 1 node failure while maintaining QUORUM. Production standard!
Use QUORUM Consistency
CL=QUORUM ensures reading latest writes while tolerating failures. Best balance!
Run Weekly Repairs
nodetool repair -pr weekly per node catches missed hints and ensures full consistency!
Monitor Hint Buildup
Watch for hint accumulation - indicates prolonged failures. Alert on excessive hints!
Smart Client Drivers
Token-aware routing automatically avoids failed nodes and retries on healthy replicas!
Multi-DC Deployment
3 datacenters with RF=3 each allows entire DC failure with zero data loss!
💼 Top 8 Interview Questions
Answer: Cassandra uses the Phi Accrual Failure Detector algorithm:
- Heartbeats: Nodes exchange heartbeats via gossip every 1 second
- Phi (Φ) Calculation: Statistical suspicion level based on heartbeat intervals
- Threshold: Φ > 8 (default) = node marked FAILED
- Detection Time: 2-3 seconds typical
- Propagation: Failure state spreads via gossip in 3-5 seconds
Advantages: Adaptive, reduces false positives, works across different networks!
Answer: Hinted handoff ensures zero data loss during temporary node failures:
- Mechanism: Coordinator stores missed writes ("hints") for down replicas
- Storage: Hints saved in system.hints table on coordinator
- TTL: 3 hours default (max_hint_window_in_ms)
- Replay: When node returns, coordinator automatically replays hints
- Result: All replicas synchronized without manual intervention
Limitation: Outages > 3 hours require nodetool repair!
Answer: Read repair automatically fixes inconsistencies during reads:
- Process: Coordinator queries all replicas, compares timestamps
- Detection: Identifies stale replicas with older timestamps
- Action: Returns latest to client, sends latest to stale replicas in background
- Types:
- Foreground: Automatic with CL > ONE
- Background: Periodic random reads (configurable %)
Benefit: Transparent consistency maintenance without manual intervention!
Answer:
Writes:
- Sent to remaining healthy replicas
- If enough replicas respond (QUORUM), write succeeds
- Coordinator stores hint for failed replica
- Client sees no error if CL requirements met
Reads:
- Client driver automatically reroutes to healthy replicas
- If enough replicas respond (QUORUM), read succeeds
- Read repair fixes any inconsistencies
- Zero impact with RF=3, CL=QUORUM
Result: Zero downtime for properly configured clusters!
Answer: Run nodetool repair in these scenarios:
- After Long Outages: Node down > 3 hours (hints expired)
- Regular Maintenance: Weekly repair per node (best practice)
- After Data Changes: Node additions, removals, or datacenter changes
- Suspected Inconsistency: Application seeing stale reads
- Before gc_grace_seconds: Ensure deletes propagated (default 10 days)
Command: nodetool repair -pr (primary range only)
Duration: Minutes to hours depending on data size
Answer: Replication Factor (RF) directly determines failure tolerance:
- RF=1: Zero tolerance (any failure = data loss)
- RF=2: Minimal (can't tolerate failures with QUORUM)
- RF=3: Tolerates 1 node failure (QUORUM=2, 2 nodes survive)
- RF=5: Tolerates 2 node failures (QUORUM=3, 3 nodes survive)
Formula: Max failures = RF - 1 (with QUORUM)
Trade-off: Higher RF = better availability but more storage cost
Recommendation: RF=3 minimum for production!
Answer: Losing quorum means operations fail:
Example (RF=3, CL=QUORUM):
- Normal: 3 nodes → QUORUM=2 ✅ Works
- 1 failure: 2 nodes → QUORUM=2 ✅ Still works
- 2 failures: 1 node → QUORUM=2 ❌ FAILED!
Impact:
- UnavailableException thrown
- Reads fail
- Writes rejected
- Application must handle errors
Solutions:
- Emergency: Temporarily lower CL to ONE
- Proper: Restore failed nodes quickly
- Prevention: Use RF=5 or multi-DC deployment
Answer: Automatic recovery process:
- T=0-10s: Node starts, contacts seeds, rejoins gossip
- T=10-30s: Gossip propagates "node UP" cluster-wide
- T=30s-10min: Coordinators detect node, replay stored hints
- T=5-10min: All hints replayed, node fully operational
Hint Replay:
- Coordinators send missed writes with original timestamps
- Node applies writes, catches up
- Hints deleted after successful replay
- Completely automatic!
Special Cases: If down > 3 hours, run nodetool repair!
Responsive Ad