Section 2: Core Architecture

Node Failure Handling

Master Cassandra's fault tolerance with animated diagrams showing failure detection, hinted handoff, and automatic recovery!

Advertisement

📖 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

Phi Accrual Failure Detection (Animated) Node A HEALTHY 💚 Heartbeat Node B Monitoring T=0s Φ = 0 T=1s Φ = 3 T=2s Φ = 6 T=3s Φ > 8 ❌ 🚨 Node Marked as FAILED 1. Gossip spreads failure state (3-5 seconds) 2. Client drivers stop routing to failed node 3. Requests redirected to healthy replicas 4. Hinted handoff activated for writes 5. Zero user-visible impact! ✅

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

Complete Failure Handling Process BEFORE: All Nodes Healthy N1 N2 N3 FAILURE: Node 2 Goes Down N1 Detects N2 DOWN N3 Detects IMMEDIATE ACTIONS (0-5 seconds) 1. 🗣️ Gossip Update - Mark N2 as DOWN 2. 🔀 Reroute Requests - Stop sending to N2 3. 📬 Hinted Handoff - N1 & N3 store hints 4. ✅ Zero Impact - Cluster continues normally! RECOVERY: Node 2 Returns N1 N2 N3 Replay hints Replay hints ✅ Fully synchronized - back to normal!

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

Hinted Handoff Mechanism (Animated) 👤 Client WRITE Coord Coordinator 1. Write R1 ✅ Success R2 ❌ DOWN R3 ✅ Success 📬 HINT Coordinator stores: For: Replica 2 Data: user123=value Timestamp: 10:30:00 TTL: 3 hours How It Works 1. Write to 3 replicas 2. R2 is down - write fails 3. R1 & R3 success = QUORUM ✅ 4. Coordinator stores hint 5. When R2 returns, replay hint 6. All replicas synchronized!

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 Process (Animated) 👤 Client READ Coord Replica 1 user123 = "New Value" TS: 10:30:00 ✅ LATEST Replica 2 user123 = "Old Value" TS: 09:00:00 ⚠️ STALE Replica 3 user123 = "Very Old" TS: 08:00:00 ⚠️ VERY STALE Coordinator Actions 1. Compare timestamps → R1 has latest (10:30:00) 2. Return latest value to client ("New Value") ✅ 3. Background: Send latest to R2 & R3 → All synced! 🔧

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

1️⃣

Set RF=3 Minimum

Replication Factor 3 allows tolerating 1 node failure while maintaining QUORUM. Production standard!

2️⃣

Use QUORUM Consistency

CL=QUORUM ensures reading latest writes while tolerating failures. Best balance!

3️⃣

Run Weekly Repairs

nodetool repair -pr weekly per node catches missed hints and ensures full consistency!

4️⃣

Monitor Hint Buildup

Watch for hint accumulation - indicates prolonged failures. Alert on excessive hints!

5️⃣

Smart Client Drivers

Token-aware routing automatically avoids failed nodes and retries on healthy replicas!

6️⃣

Multi-DC Deployment

3 datacenters with RF=3 each allows entire DC failure with zero data loss!

💼 Top 8 Interview Questions

1
How does Cassandra detect node failures?
+

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!

2
What is hinted handoff?
+

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!

3
What is read 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!

4
What happens to reads/writes when a node fails?
+

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!

5
When should you run nodetool repair?
+

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

6
How does RF affect fault tolerance?
+

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!

7
What happens if you lose quorum?
+

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
8
What happens when failed node returns?
+

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!

Advertisement

Responsive Ad