Cassandra Hinted Handoff
Master write availability during failures! Deep dive into hint storage, replay mechanisms, TTLs, and the automatic recovery process that ensures eventual consistency even when nodes are down.
📖 The Story: The Unreliable Delivery Service
Imagine you're sending packages through ShipFast Delivery - but recipients aren't always home. How do you ensure packages get delivered eventually?
❌ THE BAD APPROACH: Give Up Immediately
The Problem: Driver rings doorbell → nobody home → package returned to sender
├─ Attempt 1 (Monday 10am) → No answer ❌
├─ Return to sender
└─ Alice never gets her package! 😡
📦 Package for Bob at 456 Oak Ave:
├─ Attempt 1 (Tuesday 2pm) → No answer ❌
├─ Return to sender
└─ Bob never gets his package! 😡
Result:
• Customers angry (didn't get packages)
• Senders frustrated (wasted shipping)
• Business loses customers!
✅ THE BRILLIANT APPROACH: Leave a Note & Retry
The Solution: Driver can't deliver → leaves NOTE at nearest postal depot → retries later
Attempt 1 (Monday 10am):
├─ Ring doorbell → No answer ❌
├─ Driver writes NOTE: "Package for 123 Main St"
├─ Leaves note at Depot #5 (nearest)
└─ Continues deliveries to other addresses ✅
Alice returns home (Monday 6pm):
├─ Depot #5 sees Alice is back (lights on!)
├─ Driver reads NOTE: "Package for 123 Main St"
├─ Delivers package successfully! 📦 → 🏠
└─ Alice gets her package! ✅
TTL (Time To Live) = 3 hours:
If Alice doesn't return within 3 hours:
├─ NOTE expires and is thrown away
├─ Package returns to sender
└─ Sender can resend or contact Alice directly
Why It's BRILLIANT:
- Temporary Failure = OK: Driver doesn't give up after one attempt
- Automatic Retry: When Alice returns, package delivered automatically
- No Blocking: Driver continues other deliveries (doesn't wait)
- Depot Nearby: Notes stored at NEAREST depot (low overhead)
- TTL Protection: Notes don't pile up forever (3-hour limit)
- Eventually Delivered: Alice gets package when she returns!
🎯 This is EXACTLY Cassandra's Hinted Handoff!
- Package = Write (INSERT/UPDATE)
- Alice's House = Target Replica Node
- Nobody Home = Node Down/Unreachable
- NOTE = Hint (stored mutation)
- Postal Depot = Coordinator Node
- Alice Returns = Node Comes Back Online
- Delivery Retry = Hint Replay
- 3-Hour TTL = max_hint_window (default 3 hours)
Client writes data (RF=3):
INSERT INTO users (id, name) VALUES ('alice', 'Alice Smith');
Coordinator identifies replicas:
Node 5 (PRIMARY), Node 12, Node 19
Problem: Node 19 is down! ❌
├─ Node 5: Write successful ✅
├─ Node 12: Write successful ✅
├─ Node 19: TIMEOUT ❌ (down for maintenance)
└─ Coordinator stores HINT for Node 19
Hint stored on Coordinator:
"When Node 19 returns, write alice='Alice Smith'"
Node 19 returns (30 minutes later):
├─ Coordinator detects Node 19 is UP
├─ Reads HINT: "Write alice='Alice Smith'"
├─ Replays write to Node 19 ✅
└─ All 3 replicas now have data! ✅
Result: Write succeeded despite Node 19 failure!
Eventual consistency achieved automatically!
ShipFast's delivery notes = Cassandra's Hinted Handoff!
Same brilliance, zero data loss!
💾 What is Hinted Handoff?
The mechanism that ensures writes succeed even when replica nodes are temporarily unavailable.
Complete Definition
Hinted Handoff: When a replica node is unavailable during a write operation, the coordinator node temporarily stores the write (a "hint") and automatically replays it to the target node when it comes back online, ensuring eventual consistency without client intervention.
Key Components:
- Hint: A stored mutation (write) meant for an unavailable node
- Coordinator: Node that stores hints for failed replicas
- Target Node: The unavailable replica that should have received the write
- Hint Directory: On-disk storage for hints (separate from data)
- Replay: Automatic delivery of hints when target node recovers
- TTL: Time limit for hints (default 3 hours)
Why Hinted Handoff is Critical:
- Write Availability: Writes succeed even when replicas are down
- No Client Retry: Automatic recovery without application changes
- Eventual Consistency: Data eventually reaches all replicas
- Performance: Faster writes (don't wait for slow/down nodes)
- Operational Simplicity: Node maintenance without write failures
Without Hinted Handoff
Traditional databases
Write arrives:
├─ Node 1: Success ✅
├─ Node 2: Success ✅
├─ Node 3: DOWN ❌
└─ QUORUM = 2/3 achieved
Problem:
Node 3 NEVER gets the data!
Data inconsistent until manual repair!
- Consistency: Permanent inconsistency
- Recovery: Manual repair required
- Operations: Complex
With Hinted Handoff
Cassandra approach
Write arrives:
├─ Node 1: Success ✅
├─ Node 2: Success ✅
├─ Node 3: DOWN ❌
├─ QUORUM = 2/3 achieved
└─ Coordinator stores HINT
Node 3 returns (later):
├─ Coordinator detects UP
├─ Replays hint automatically
└─ Node 3 has data! ✅
- Consistency: Automatic eventual consistency
- Recovery: Zero manual intervention
- Operations: Simple and automatic
⚙️ How Hinted Handoff Works: Step-by-Step
Complete flow from write failure through hint storage to automatic replay.
Expert: Detailed Process Breakdown
Step 1: Write Request
- Client sends write to any node (becomes coordinator)
- Coordinator calculates partition key hash → token
- Determines replica nodes using token ring (e.g., Node 5, 12, 19)
- Attempts to write to all RF replicas in parallel
Step 2: Failure Detection
- Write to Node 5: Success (50ms) ✅
- Write to Node 12: Success (45ms) ✅
- Write to Node 19: TIMEOUT after 2 seconds ❌
- QUORUM (2/3) achieved → Write successful from client perspective
Step 3: Hint Creation
- Coordinator creates hint: Serialized mutation + target node ID + timestamp
- Hint stored in: /var/lib/cassandra/hints/hint_{target_node_id}_{timestamp}.hints
- Hint includes: Full mutation, TTL (3 hours), target node UUID
- Coordinator continues processing other requests (non-blocking)
Step 4: Gossip Detection
- Coordinator gossips with other nodes every 1 second
- When Node 19 returns: Gossip discovers it's UP
- Coordinator checks: "Do I have hints for Node 19?" → YES!
Step 5: Hint Replay
- Coordinator reads all hints for Node 19 from disk
- Sends hints to Node 19 (batch of mutations)
- Node 19 writes data to memtable and commit log
- On success: Coordinator deletes hint file
- On failure: Keeps hint, retries later
Step 6: Cleanup
- Successful replay: Delete hint file immediately
- TTL expired (>3 hours): Delete hint without replay
- Disk space: Hints use separate directory, limited to max_hint_window
📦 Hint Storage: Under the Hood
How hints are stored on disk and managed over time.
Storage Architecture
Hint File Location:
├─ hint_node19_1701234567890.hints (for Node 19)
├─ hint_node23_1701234589012.hints (for Node 23)
└─ hint_node45_1701234601234.hints (for Node 45)
Hint File Structure:
- Target Node UUID: Which node should receive this hint
- Mutation: Complete serialized write (partition key, clustering, columns, timestamp)
- Creation Time: When hint was created (for TTL calculation)
- Keyspace/Table: Where to apply the mutation
- Hint ID: Unique identifier for deduplication
Storage Characteristics:
- Separate Directory: Hints isolated from data files
- Sequential Writes: Fast append-only writes
- Per-Node Files: One file per target node
- Size Limits: Configurable max_hints_file_size (128MB default)
- Compression: Optionally compressed (LZ4)
Hint File Format
Magic bytes: 0xC5A55A1D
Version: 2
Target node: a1b2c3d4-...
Hint Entry 1:
Timestamp: 1701234567890
Keyspace: users
Table: user_profiles
Partition: alice
Mutation: [serialized]
Size: 256 bytes
Hint Entry 2:
Timestamp: 1701234567891
Keyspace: users
Table: user_sessions
Partition: bob
Mutation: [serialized]
Size: 312 bytes
Footer:
Checksum: CRC32
Total hints: 2
Storage Management
- Creation: Hint file created when first write to unavailable node fails
- Append: Additional hints appended to same file for same target
- Rotation: New file when max_hints_file_size reached
- Replay: Read sequentially, batch send to target
- Deletion: Deleted after successful replay or TTL expiry
- Cleanup: Background thread removes expired hints
Performance Impact:
- Negligible write overhead (~100μs per hint)
- No impact on normal writes
- Async replay (doesn't block queries)
- Disk I/O: Sequential (fast!)
Storage Considerations
Disk Space Management:
- Hints Directory Size: Can grow if nodes stay down for extended period
- Monitoring: Watch disk usage in hints directory
- TTL Protection: max_hint_window (3 hours) prevents unbounded growth
- Manual Cleanup: Can delete hint files if target node permanently lost
Example Disk Usage:
Total writes missed: 1000 * 60 * 60 * 2 = 7.2M writes
Average hint size: 512 bytes
Total hints size: 7.2M * 512 bytes = 3.6GB
With compression (LZ4): ~1.8GB
Replay time (10K hints/sec): ~12 minutes
⚙️ Configuration & Tuning
Complete guide to hinted handoff configuration parameters.
Configuration in cassandra.yaml
hinted_handoff_enabled: true
max_hint_window_in_ms: 10800000 # 3 hours
max_hints_delivery_threads: 2
max_hints_file_size_in_mb: 128
hints_flush_period_in_ms: 10000 # 10 seconds
hints_directory: /var/lib/cassandra/hints
hints_compression:
class_name: LZ4Compressor
# Advanced Settings
max_hints_size_per_host_in_mb: 4096 # 4GB per node
hint_window_persistent_enabled: true
batchlog_replay_throttle_in_kb: 1024
Production Recommendations:
- Always Enabled: Keep hinted_handoff_enabled: true (critical for availability)
- TTL Window: 3 hours (default) works for 99% of cases, increase to 6-12 hours for planned maintenance
- Delivery Threads: 2-4 threads sufficient for most clusters, 8+ for very large clusters (500+ nodes)
- File Size: 128MB fine for normal workloads, increase to 256-512MB for high-write systems
- Compression: LZ4 (best balance of speed and compression)
🖥️ Interactive Hinted Handoff Simulator
Visualize how hints are created, stored, and replayed in real-time!
Simulate writes with node failures and watch hints being created and replayed.
Try this sequence:
1. Click "Simulate Write" (all nodes UP) → Normal replication
2. Click "Node Goes Down" → Node 3 marked as DOWN
3. Click "Simulate Write" → Hint created for Node 3
4. Click "Node Comes Up" → Hints automatically replayed!
RF=3, Nodes: 1, 2, 3
✅ Best Practices & Common Pitfalls
Production-grade recommendations for hinted handoff.
Best Practices
- Keep Enabled: Never disable hinted_handoff_enabled in production
- Monitor Hints: Track hint queue size (nodetool hints)
- Disk Space: Ensure 20-30% free space in hints directory
- TTL Window: Match max_hint_window to typical downtime + buffer
- Repair Regularly: Run nodetool repair weekly (hints ≠ repair)
- Graceful Maintenance: Use nodetool drain before shutting down nodes
- Alert on Buildup: Alert when hints directory >5GB
- Test Recovery: Simulate node failures in staging
Common Pitfalls
- Disabling Handoff: Turning off hints = lost writes during failures
- Ignoring Disk Space: Hints directory fills up = new hints rejected
- Relying Only on Hints: Hints expire after TTL, still need repair
- Long Downtime: Node down >3 hours = hints lost, manual repair needed
- No Monitoring: Hint buildup goes unnoticed until disk full
- Wrong TTL: Too short = data loss, too long = disk waste
- Insufficient Threads: Slow replay = hints accumulate
- No Testing: First failure in production reveals issues
Monitoring Hinted Handoff
Key Metrics to Track:
- Hints Stored: Total number of hints currently on disk
- Hints in Flight: Hints being replayed right now
- Hints Dropped: Hints expired/discarded (should be zero!)
- Replay Latency: Time to replay hints to recovered node
- Disk Usage: Size of hints directory
Useful Commands:
nodetool hintstats
# View hints for specific node
nodetool getsstables keyspace table key
# Manually trigger hint replay
nodetool pausehandoff
nodetool resumehandoff
# Check disk usage
du -sh /var/lib/cassandra/hints/
💼 Interview Questions & Answers
Master these 25 essential questions on Cassandra Hinted Handoff - from beginner to expert level!
Answer:
Hinted Handoff is a mechanism that ensures writes succeed even when replica nodes are temporarily unavailable by storing "hints" (failed writes) and automatically replaying them when the node recovers.
How It Works:
- Write arrives for RF=3 nodes (N1, N2, N3)
- N3 is down/unreachable
- Coordinator stores hint: "When N3 returns, apply this write"
- Write succeeds with QUORUM (2/3)
- When N3 comes back, coordinator automatically replays hints
Why Critical:
- Write Availability: Writes don't fail due to temporary node issues
- Automatic Recovery: No manual intervention required
- Eventual Consistency: All replicas eventually get the data
- Operational Simplicity: Can perform maintenance without impacting writes
Answer:
Default TTL: 3 hours (10,800,000 milliseconds)
Parameter: max_hint_window_in_ms in cassandra.yaml
Why TTL Exists:
- Disk Space Protection: Prevents hints from filling up disk if node stays down for extended period
- Bounded Storage: Limits maximum hints storage per node
- Reasonable Window: 3 hours covers typical maintenance windows and transient failures
- Forces Repair: After 3 hours, you should run nodetool repair instead of relying on hints
What Happens After TTL:
- Hints older than 3 hours are automatically deleted
- Down node won't receive those writes via hints
- Must run nodetool repair to synchronize data
- Read repair or anti-entropy repair will eventually fix inconsistencies
When to Increase TTL:
- Planned maintenance longer than 3 hours
- Slow node recovery expected
- Delayed hardware replacement scenarios
- Typical: Increase to 6-12 hours for planned downtime
Answer:
Location: /var/lib/cassandra/hints/ (configurable via hints_directory)
File Structure:
├─ hint_a1b2c3d4-..._1701234567890.hints
├─ hint_e5f6g7h8-..._1701234589012.hints
└─ hint_i9j0k1l2-..._1701234601234.hints
Naming: hint_{target_node_uuid}_{timestamp}.hints
Storage Characteristics:
- Separate Directory: Isolated from data files and commit logs
- Per-Node Files: One or more files per target node
- Sequential Writes: Append-only for fast writes
- File Rotation: New file when max_hints_file_size_in_mb (128MB default) reached
- Compression: LZ4 compression by default (configurable)
Important Notes:
- Hints are stored on the COORDINATOR node, not on the target node
- Each coordinator stores hints for nodes it tried to write to
- Multiple coordinators may have hints for the same target node
- Deleted after successful replay or TTL expiration
Answer:
Cassandra uses the Gossip protocol to detect when down nodes come back online and automatically triggers hint replay.
Detection Process:
- Continuous Gossip: Every node gossips with 1-3 random nodes every second
- State Updates: Nodes share state information (UP/DOWN) via gossip
- Recovery Detection: When a down node comes back up, gossip spreads this information
- Coordinator Check: Each coordinator checks "Do I have hints for this node?"
- Automatic Replay: If yes, coordinator automatically starts replaying hints
Replay Triggers:
- Node Startup: When node restarts and rejoins cluster
- Network Recovery: When network partition heals
- Periodic Check: Background thread checks every 10 seconds (hints_flush_period_in_ms)
- Manual: Can manually trigger with nodetool resumehandoff
Replay Behavior:
- Hints replayed in batches (efficient bulk transfer)
- Throttled to avoid overwhelming target node
- max_hints_delivery_threads (default 2) controls parallelism
- Successful replay → delete hint file
- Failed replay → retry later (with exponential backoff)
Answer:
NO! Hints are NOT a replacement for repair.
Why Hints ≠ Repair:
- TTL Limitation: Hints expire after 3 hours (default), repair has no time limit
- Single Coordinator: Only one coordinator stores hints for each write, if that coordinator fails, hints lost
- No Tombstone Propagation: Deleted data (tombstones) aren't hinted
- Partial Coverage: Hints only cover writes during downtime, not pre-existing inconsistencies
- Disk Space: Coordinators may drop hints if disk space exhausted
What Hints Do Well:
- Handle transient failures (node restarts, brief network issues)
- Maintain write availability during short downtimes
- Reduce immediate inconsistency window
- Speed up convergence for recent writes
When You MUST Run Repair:
- Node down > max_hint_window (3 hours)
- Node permanently failed and replaced
- Data corruption detected
- Regular maintenance (weekly recommended)
- After major network partition
Best Practice:
Run nodetool repair weekly on all nodes regardless of hints. Treat hints as an optimization, not a guarantee.
Answer:
The hints stored on that coordinator are lost! This is a fundamental limitation of hinted handoff.
Scenario:
- Node A (coordinator) writes to nodes B, C, D (RF=3)
- Node D is down → Node A stores hint for D
- Node A crashes before D recovers
- Hint is lost → Node D never receives this write via hints
Why This Is OK:
- QUORUM Still Met: Write succeeded to B and C (2/3)
- Read Repair: Next read will trigger read repair to D
- Anti-Entropy Repair: Scheduled repair will sync D
- Multiple Coordinators: Different writes use different coordinators (distributed risk)
Mitigation:
- Use RF=3 minimum (more replicas = more redundancy)
- Run regular repairs (weekly)
- Monitor coordinator health
- Use consistency level QUORUM or higher
Answer:
Minimal impact! Hints are designed to be lightweight and non-blocking.
Performance Characteristics:
- Write Overhead: ~100 microseconds per hint (negligible)
- Non-Blocking: Hint storage is asynchronous (doesn't delay client response)
- Sequential I/O: Hints written sequentially (fast disk operations)
- Memory Usage: Small buffer in memory, flushed every 10 seconds
- No Query Impact: Hint storage doesn't affect read/write path
Write Flow with Hints:
T=1ms: Coordinator identifies replicas
T=2ms: Parallel writes to N1, N2, N3
T=5ms: N1 responds OK
T=6ms: N2 responds OK
T=2002ms: N3 timeout (down)
T=2002.1ms: Create hint (0.1ms)
T=6ms: Client receives SUCCESS (QUORUM achieved)
Note: Client doesn't wait for hint creation!
Replay Performance:
- Replay is throttled (doesn't overwhelm target)
- Batched transfers (efficient network usage)
- Background operation (doesn't impact foreground traffic)
- Typical: 10K-50K hints/second replay rate
Answer: No, hinted handoff is a global cluster-level setting (hinted_handoff_enabled) and cannot be disabled per keyspace or table. It's all-or-nothing for the entire cluster. You can disable it globally but this is NOT recommended for production as it severely impacts write availability during node failures.
Answer: max_hints_delivery_threads (default 2) controls how many threads are dedicated to replaying hints. Increase to 4-8 for large clusters (500+ nodes) or when you need faster hint replay after prolonged downtime. More threads = faster replay but more CPU/network usage. Monitor CPU during replay and adjust accordingly.
Answer: Hints are created regardless of consistency level. With LOCAL_QUORUM (RF=3 in DC), if 1 local replica is down, hint created and write succeeds (2/3 local replicas). The hint ensures the down node eventually gets the data. Consistency level only determines how many acks needed for write success, but hints are always created for any failed replica write attempts.
Answer: Use: (1) nodetool hintstats - shows hints per node, (2) du -sh /var/lib/cassandra/hints/ - disk usage, (3) JMX metrics: org.apache.cassandra.metrics.StorageMetrics.TotalHints, (4) Set alerts when hints directory >5GB or hints age >30 minutes.
Answer: New hints are dropped (not stored). Writes still succeed if QUORUM met, but down node won't receive data via hints. Must run repair to sync. Set max_hints_size_per_host_in_mb to prevent unbounded growth. Monitor disk usage and alert at 80% capacity.
Answer: Yes, by default hints are compressed using LZ4 (hints_compression parameter). LZ4 provides good compression (~50-70%) with minimal CPU overhead. Can also use Snappy or disable compression. Compression reduces disk usage significantly for text-heavy data.
Answer: Yes, you can manually delete hint files from /var/lib/cassandra/hints/ if target node permanently lost or to free space. However, you MUST run full repair afterward to ensure data consistency. Use nodetool truncatehints (Cassandra 4.0+) for safe deletion or stop Cassandra and rm hint files.
Answer: By default, hints only stored for nodes in SAME datacenter (hinted_handoff_enabled_by_dc). Cross-DC hints disabled to save WAN bandwidth. If remote DC node down, use repair instead of hints. Can enable cross-DC hints but not recommended due to network costs and latency.
Answer: Controls how often hints are flushed from memory to disk (default 10 seconds). Lower value (5s) = faster recovery potential but more I/O. Higher value (30s) = less I/O but hints only persisted every 30s. If coordinator crashes before flush, recent hints lost. 10s is good balance.
Answer: Yes! Use nodetool pausehandoff to pause hint replay (useful during maintenance) and nodetool resumehandoff to restart. Pausing doesn't prevent new hints from being created, only stops replay to recovered nodes. Use during high-traffic periods to avoid replay overhead.
Answer: Depends on hint volume and network speed. Typical: 10K-50K hints/second. Example: Node down 1 hour with 1000 writes/sec = 3.6M hints. At 30K hints/sec = 2 minutes replay time. Large backlogs (node down 3 hours, high write rate) might take 10-30 minutes. Monitor with nodetool hintstats.
Answer: nodetool drain flushes memtables, stops accepting writes, but hints still created by other nodes for this node during downtime. Drain doesn't prevent hint creation - it just ensures clean shutdown. When node restarts, other coordinators replay their hints to it. Always use drain before planned restarts.
Answer: Yes! Tombstones (DELETE mutations) are treated like any other write. If node down during DELETE, hint created with tombstone. This ensures deleted data doesn't "resurrect" on recovered node. Critical for data correctness. Tombstones must reach all replicas before gc_grace_seconds expires.
Answer: Each mutation in a batch can create separate hints. If batch has 10 mutations and 1 replica is down, 10 separate hints created. Batches themselves aren't stored as atomic units in hints. When replayed, hints applied individually. Batch semantics (logged/unlogged) don't affect hint creation.
Answer: Hints handle individual failed replica writes. Batchlog handles failed LOGGED BATCH statements (ensures atomicity across partitions). Hints: temporary storage for failed writes to replicas. Batchlog: temporary storage for entire batches during coordinator failures. Both provide eventual consistency but for different failure scenarios.
Answer: No! Cassandra uses timestamps for conflict resolution (Last Write Wins). Even if hints arrive out of order, the timestamp determines which value wins. Hint with older timestamp doesn't overwrite newer data. This is why client-side timestamps and NTP synchronization are critical.
Answer: Check: (1) system.log for hint replay errors, (2) nodetool hintstats for replay progress, (3) Verify target node is UP via nodetool status, (4) Check network connectivity, (5) Verify disk space on both coordinator and target, (6) Review cassandra.yaml settings (max_hints_delivery_threads), (7) Use nodetool pausehandoff / resumehandoff to control replay.
Answer:
Comprehensive Strategy:
- Multi-DC Setup: RF=3 per DC, 3 DCs minimum (9 total copies)
- Consistency: Write EACH_QUORUM (ensures QUORUM in every DC), Read LOCAL_QUORUM
- Hint Configuration: max_hint_window=21600000 (6 hours) for extended maintenance windows
- Repair Schedule: Weekly full repair per DC (don't rely on hints alone)
- Monitoring: Alert when hints directory >10GB or hints age >1 hour
- Failover: If DC down >6 hours, run immediate repair when recovered (don't wait for weekly)
- Backup: Daily snapshots independent of hints (hints not backups!)
- Testing: Quarterly DR drills simulating DC failure and recovery
Key Principle: Hints are optimization for short-term failures. Repair is guarantee for long-term consistency.
🎓 Chapter Summary
You now understand Cassandra's Hinted Handoff mechanism!
Key Concepts:
- Hints ensure writes succeed even when replicas temporarily unavailable
- Stored on coordinator node, replayed automatically when target recovers
- Default 3-hour TTL prevents unbounded disk growth
- Minimal performance impact, high availability benefit
- NOT a replacement for repair (complementary mechanism)
- Multi-DC deployments use same-DC hints only by default
Production Best Practices:
- Never disable hinted_handoff_enabled in production
- Run weekly repairs regardless of hints
- Monitor hints directory size (alert at 5-10GB)
- Increase max_hint_window for planned maintenance (6-12 hours)
- Use graceful shutdown (nodetool drain) before restarts
Responsive Ad