Replication Factor
Master data redundancy! Learn how Cassandra replicates data across nodes for fault tolerance, high availability, and zero downtime!
📖 Netflix's Data Replication: 99.99% Availability
Netflix stores viewing history across 3 data centers with RF=3 in each. Challenge: Ensure zero data loss even if entire datacenter fails. Solution: RF=3 per DC + 3 DCs = 9 total copies! Result: When AWS US-East-1 went down in 2015, Netflix continued streaming seamlessly using other datacenters. 100 million viewers saw zero disruption!
💻 Interactive Replication Console
Try RF commands in our live interactive console! Experiment with different replication configurations.
Quick Examples:
▶ Learn about RF configurations and multi-DC replication.
▶ Tip: Press Ctrl+Enter to execute commands quickly!
Learning Tips
- Experiment: Try different RF values (1, 2, 3, 5)
- Multi-DC: See how NetworkTopologyStrategy works
- Endpoints: Check which nodes store your data
- Consistency: Understanding RF + CL relationship
🎯 What is Replication Factor?
Replication Factor (RF) is the number of copies Cassandra maintains for each piece of data across the cluster.
RF = 1
Copies: 1 (no redundancy)
Fault Tolerance: None
Availability: Low
Use Case: Testing only
Risk: Any node failure = data loss
RF = 2
Copies: 2 (minimal redundancy)
Fault Tolerance: Limited
Availability: Medium
Use Case: Development
Risk: Can't use QUORUM safely
RF = 3 ⭐
Copies: 3 (standard)
Fault Tolerance: 1 node
Availability: High
Use Case: Production (recommended)
Benefit: QUORUM works perfectly
Key Concept
RF determines fault tolerance: With RF=3, cluster can tolerate (RF-1) = 2 node failures while maintaining QUORUM. This is why RF=3 is the production standard!
⚙️ How Replication Works
Cassandra replicates data to multiple nodes based on the configured RF. The coordinator ensures all replicas are updated!
# Configure keyspace with RF=3
CREATE KEYSPACE demo
WITH replication = {
'class': 'SimpleReplicationStrategy',
'replication_factor': 3
};
# Write process:
1. Client sends write to any node (becomes coordinator)
2. Coordinator calculates token: hash(user_id=123) = -8234719023847192847
3. Coordinator determines 3 replica nodes using token
4. Coordinator writes to all 3 replicas in parallel:
- Local write (if coordinator is replica): 5ms
- Network writes to other 2 replicas: 12-15ms
5. Coordinator waits for CL responses (e.g., QUORUM=2)
6. Returns success to client
Total latency: ~15ms (limited by slowest replica needed for CL)
📊 RF Scenarios Compared
Why Not RF=2?
With RF=2, QUORUM=2. This means both nodes must be available for QUORUM reads/writes. If either fails, you lose QUORUM! This is why RF=2 is NOT recommended for production despite having redundancy.
🔗 RF and Consistency Levels
RF and Consistency Level (CL) work together to determine fault tolerance and data consistency!
# Examples with RF=3 # CL=ONE (Fastest, Highest Availability) SELECT * FROM users WHERE user_id=123 USING CONSISTENCY ONE; -- Returns as soon as ANY 1 replica responds -- Can tolerate 2 node failures -- May return slightly stale data # CL=QUORUM (Balanced - RECOMMENDED) SELECT * FROM users WHERE user_id=123 USING CONSISTENCY QUORUM; -- Waits for 2 of 3 replicas (majority) -- Can tolerate 1 node failure -- Guaranteed to return latest write (if written with QUORUM) -- Best for production! # CL=ALL (Strongest, Lowest Availability) SELECT * FROM users WHERE user_id=123 USING CONSISTENCY ALL; -- Waits for ALL 3 replicas -- Cannot tolerate ANY failure -- Guaranteed latest data -- Use only when absolutely necessary # Strong Consistency Example -- Write with QUORUM (W=2) + Read with QUORUM (R=2) -- R + W = 4 > RF (3) ✓ -- This guarantees you read what you just wrote! INSERT INTO users (user_id, name) VALUES (123, 'Alice') USING CONSISTENCY QUORUM; -- W=2 SELECT * FROM users WHERE user_id=123 USING CONSISTENCY QUORUM; -- R=2 -- Guaranteed to see "Alice" (at least 1 replica overlap)
Production Best Practice
RF=3 with CL=QUORUM is the gold standard for production:
• Tolerates 1 node failure
• Balanced latency
• Strong consistency
• Used by Netflix, Instagram, Apple
🌍 Multi-Datacenter Replication
NetworkTopologyStrategy allows configuring different RF values per datacenter for geographic redundancy!
Multi-DC Consistency Levels
# LOCAL_QUORUM (Recommended for Multi-DC) # Wait for quorum in LOCAL datacenter only SELECT * FROM users WHERE user_id=123 USING CONSISTENCY LOCAL_QUORUM; -- With RF=3 per DC: Waits for 2 nodes in local DC -- Latency: ~15ms (no cross-DC wait) -- Use: Default for multi-DC setups # EACH_QUORUM (Strong Multi-DC Consistency) # Wait for quorum in EACH datacenter INSERT INTO users (user_id, name) VALUES (123, 'Alice') USING CONSISTENCY EACH_QUORUM; -- With 3 DCs, RF=3 each: Waits for 2 nodes in all 3 DCs -- Latency: ~100ms+ (cross-datacenter) -- Use: Critical data needing global consistency # LOCAL_ONE (Fastest Multi-DC) SELECT * FROM users WHERE user_id=123 USING CONSISTENCY LOCAL_ONE; -- Wait for 1 node in local DC only -- Latency: ~5ms -- Use: High availability, eventual consistency OK # Real example: Netflix Multi-DC setup # Write: LOCAL_QUORUM (fast, available) # Read: LOCAL_QUORUM (fast, consistent within DC) # Result: 99.99% availability, <20ms latency globally
Real Scenario: AWS Region Failure
In 2017, AWS S3 went down in US-East-1. Companies using Cassandra with multi-DC replication (like Netflix) stayed online by routing traffic to US-West and EU datacenters. Their RF=3 per DC meant zero data loss!
✅ Replication Best Practices
Always Use RF=3 Minimum
Production Standard: RF=3
Fault Tolerance: 1 node failure
Works With: QUORUM consistency
Used By: Netflix, Instagram, Apple
Use NetworkTopologyStrategy
For: All production clusters
Benefit: DC-aware replication
Allows: Different RF per DC
Future-proof: Easy to add DCs
Pair RF with Correct CL
RF=3: Use CL=QUORUM
Multi-DC: Use LOCAL_QUORUM
R+W > RF: Strong consistency
Balance: Availability vs consistency
Plan for Geographic Distribution
Multiple DCs: Survive region failures
RF per DC: 3 minimum
Racks: Distribute within DC
Result: 99.99%+ availability
Monitor Replication Health
Check: nodetool status regularly
Watch: Inconsistent replicas
Run: nodetool repair weekly
Alert: Node down > 3 hours
Consider Storage Costs
RF=3: 3x storage cost
Multi-DC: Multiply by DC count
Trade-off: Cost vs availability
Worth it: For critical data
Real Configuration: Apple iCloud
# Apple's reported Cassandra setup
# Cluster: 75,000+ nodes across multiple datacenters worldwide
CREATE KEYSPACE user_data WITH replication = {
'class': 'NetworkTopologyStrategy',
'US-East': 3,
'US-West': 3,
'EU-West': 3,
'Asia-Pacific': 3
};
# Write consistency: LOCAL_QUORUM
# Read consistency: LOCAL_QUORUM
# Result:
# - 99.999% availability (5 nines!)
# - < 10ms latency for local users
# - Survives entire region failure
# - Serves 1+ billion users globally
# Storage: RF=3 × 4 regions = 12 copies total
# Cost: High, but availability is priceless for iCloud
💼 Top 10 Interview Questions
Answer:
Replication Factor (RF) is the number of copies Cassandra maintains for each piece of data across the cluster.
Why it's important:
- Fault Tolerance: Can tolerate (RF-1) node failures
- High Availability: Data accessible even when nodes fail
- Performance: Reads can be distributed across replicas
- Disaster Recovery: Multiple copies protect against data loss
Example:
RF=3 means: - 3 copies of every row - Can tolerate 1 node failure with QUORUM - Can tolerate 2 node failures with CL=ONE - Used by Netflix, Instagram, Apple for 99.99%+ availability
Production Standard: RF=3 minimum!
Answer: RF=3 is the production standard because it provides the best balance:
Advantages of RF=3:
- Fault Tolerance: Can lose 1 node and maintain QUORUM
QUORUM = (3 / 2) + 1 = 2 3 nodes total, need 2 for QUORUM If 1 fails → 2 remain → QUORUM still works ✓
- Strong Consistency: QUORUM reads/writes overlap
Write QUORUM: 2 replicas updated Read QUORUM: 2 replicas checked Guaranteed overlap → see latest write ✓
- Reasonable Cost: 3x storage is acceptable trade-off
- Performance: Balanced latency (15-20ms typical)
Why NOT RF=1 or RF=2:
- RF=1: No fault tolerance, any failure = data loss
- RF=2: QUORUM=2 = can't tolerate ANY failure with QUORUM
Industry Standard: Used by all major tech companies!
Answer: RF and CL work together to determine fault tolerance and consistency guarantees.
Key Relationships:
| RF | CL | Nodes Needed | Fault Tolerance |
|---|---|---|---|
| 3 | ONE | 1 | 2 nodes |
| 3 | QUORUM | 2 | 1 node |
| 3 | ALL | 3 | 0 nodes |
Strong Consistency Formula:
For strong consistency: R + W > RF Example with RF=3: - Write CL=QUORUM (W=2) - Read CL=QUORUM (R=2) - R + W = 2 + 2 = 4 > RF (3) ✓ Result: Read always sees latest write (guaranteed overlap)
Best Practice: RF=3 with CL=QUORUM for production!
Answer:
NetworkTopologyStrategy is a replication strategy that's datacenter-aware, allowing different RF values per datacenter.
vs SimpleReplicationStrategy:
- SimpleReplicationStrategy: Single datacenter only, RF applies to entire cluster
- NetworkTopologyStrategy: Multi-datacenter aware, RF per datacenter
When to use NetworkTopologyStrategy:
- ✅ Always - Even for single DC! (future-proof)
- ✅ Multi-datacenter deployments
- ✅ Production environments
- ✅ When you need rack awareness
Example:
-- Single DC (still use NTS!)
CREATE KEYSPACE demo WITH replication = {
'class': 'NetworkTopologyStrategy',
'datacenter1': 3
};
-- Multi DC (different RF per DC)
CREATE KEYSPACE global WITH replication = {
'class': 'NetworkTopologyStrategy',
'US-East': 3,
'US-West': 3,
'EU-West': 2 -- Lower RF in less critical DC
};
-- Netflix example
CREATE KEYSPACE viewing_history WITH replication = {
'class': 'NetworkTopologyStrategy',
'US-East-1': 3,
'US-West-2': 3,
'EU-West-1': 3
};
-- 9 total copies (3 per DC), survives entire DC failure!
Benefits:
- Survive datacenter failures
- Optimize for geographic distribution
- Control storage costs per region
- Rack-aware placement (avoid failure zones)
Answer: Use nodetool getendpoints command:
# Syntax nodetool getendpoints
# Example: Find replicas for user_id=12345 $ nodetool getendpoints demo users 12345 10.1.0.1 10.1.0.2 10.1.0.3 # This shows the 3 nodes (RF=3) that store this user's data # For composite partition key $ nodetool getendpoints demo sensor_data 'sensor-123' '2024-12-26' 10.1.0.2 10.1.0.3 10.1.0.4 # Multi-DC example $ nodetool getendpoints global users 67890 # Output (with 2 DCs, RF=3 each): 10.1.0.1 # US-East 10.1.0.2 # US-East 10.1.0.3 # US-East 10.2.0.1 # US-West 10.2.0.2 # US-West 10.2.0.3 # US-West # Total: 6 replicas (3 per DC) # You can also check the token $ cqlsh -e "SELECT token(user_id) FROM users WHERE user_id=12345" system.token(user_id) ----------------------- -8234719023847192847 # Then check which nodes own that token range $ nodetool ring demo | grep -- "-8234719023847192847" This is useful for:
- Debugging data distribution
- Verifying RF configuration
- Understanding query routing
- Troubleshooting replication issues
6Can you change RF after keyspace is created?Answer: Yes! But requires repair to redistribute data.
Process:
# Step 1: Alter keyspace ALTER KEYSPACE demo WITH replication = { 'class': 'NetworkTopologyStrategy', 'datacenter1': 3 -- Changed from RF=2 to RF=3 }; # Step 2: Run repair on ALL nodes (very important!) # This copies data to the new replicas $ nodetool repair -full demo # Or repair each node sequentially $ for node in node1 node2 node3; do ssh $node "nodetool repair -full demo" done # Step 3: Verify replication $ nodetool getendpoints demo users 12345 # Should now show 3 nodes (was 2 before)Important Notes:
- ⚠️ Must run repair - Schema change alone doesn't move data!
- ⚠️ Repair can take hours on large clusters
- ⚠️ Run repair on ALL nodes, not just new replicas
- ✅ No downtime required
- ✅ Can do rolling repair (one node at a time)
Decreasing RF:
# Change from RF=3 to RF=2 ALTER KEYSPACE demo WITH replication = { 'class': 'NetworkTopologyStrategy', 'datacenter1': 2 }; # Data automatically removed from old replicas during cleanup $ nodetool cleanup demo # This frees up space on nodes no longer storing replicasBest Practice: Plan RF carefully upfront to avoid costly repairs!
7What happens if all replicas for a partition fail?Answer: That data becomes temporarily unavailable.
Scenario with RF=3:
Data for user_id=12345 stored on nodes: N1, N2, N3 If ALL 3 nodes fail simultaneously: - Reads: UnavailableException (no replicas to read from) - Writes: UnavailableException (no replicas to write to) - Data: Temporarily unavailable (not lost!) When any 1 node returns: - CL=ONE queries work again - Data accessible - Service restoredProbability:
- Extremely rare with proper infrastructure
- Nodes chosen from different racks/availability zones
- Simultaneous failure of 3 independent nodes: ~0.001% chance
Protection Strategies:
- Rack Awareness: Place replicas in different racks
# cassandra-rackdc.properties dc=US-East rack=rack1 # Distribute across rack1, rack2, rack3- Multi-DC: RF=3 in multiple datacenters
'US-East': 3, # Even if all 3 fail here... 'US-West': 3 # ...data still available in US-West!- Higher RF: RF=5 means need 5 failures (very unlikely)
- Monitoring: Alert on first node failure, fix quickly
Real World Example:
Netflix with Multi-DC setup: - US-East DC loses power (all 3 replicas there down) - US-West DC still has 3 replicas → Service continues - EU-West DC still has 3 replicas → Service continues - Users in US-East routed to US-West (slightly higher latency) - Zero data loss, zero downtime!8What is LOCAL_QUORUM and when should you use it?Answer:
LOCAL_QUORUM requires a quorum of replicas to respond, but only from the LOCAL datacenter (not cross-DC).
Comparison:
Consistency Level Behavior Latency QUORUM Wait for majority across ALL DCs High (~100ms+) LOCAL_QUORUM Wait for majority in LOCAL DC only Low (~15ms) EACH_QUORUM Wait for majority in EACH DC Very High (~200ms+) When to use LOCAL_QUORUM:
- ✅ Multi-DC deployments (recommended default)
- ✅ When low latency is priority
- ✅ When eventual cross-DC consistency is acceptable
- ✅ For 99% of multi-DC use cases
Example:
# Setup: 2 DCs with RF=3 each # Client in US-East datacenter SELECT * FROM users WHERE user_id=12345 USING CONSISTENCY LOCAL_QUORUM; # Behavior: - Checks only US-East replicas (N1, N2, N3) - Waits for 2 of 3 to respond (QUORUM in local DC) - Does NOT wait for US-West replicas - Latency: ~15ms (local DC only) # Contrast with QUORUM: USING CONSISTENCY QUORUM; # Would wait for 4 of 6 total replicas (across both DCs) # Latency: ~100ms (cross-DC network delay) # Write propagates to other DC asynchronously INSERT INTO users VALUES (...) USING CONSISTENCY LOCAL_QUORUM; # Written to US-East immediately (15ms) # Replicated to US-West in background (eventual)Best Practice for Multi-DC:
- Writes: LOCAL_QUORUM (fast, available)
- Reads: LOCAL_QUORUM (fast, consistent within DC)
- Critical Writes: EACH_QUORUM (slow but guaranteed globally)
6How does RF affect write performance?Answer: Higher RF increases write latency but not as much as you might expect!
Write Process:
1. Coordinator writes to ALL RF replicas in PARALLEL 2. Waits for CL responses 3. Returns success to client Key point: Writes are PARALLEL, not sequential!Latency Comparison (CL=QUORUM):
RF Replicas Written Parallel Writes Latency 1 1 N/A ~5ms 2 2 Yes ~12ms 3 3 (wait for 2) Yes ~15ms 5 5 (wait for 3) Yes ~18ms Why the small increase?
- Writes happen in parallel to all replicas
- Latency = slowest replica needed for CL
- Network is typically fast within datacenter (1-5ms)
- Most time spent in disk writes, not network
Other Performance Impacts:
- Network Bandwidth: More replicas = more network traffic
RF=3: 3x network bandwidth vs RF=1 For 1000 writes/sec with 1KB each: RF=1: 1 MB/s RF=3: 3 MB/s total (not per node)- Disk I/O: Each replica writes to disk
RF=3: Each node writes 3x less data than RF=1 (data distributed across more nodes)- Compaction: More replicas = more compaction overhead
Real-World Benchmark:
Netflix observed latencies (CL=LOCAL_QUORUM): RF=3, single DC: p50=12ms, p99=45ms RF=3, multi-DC: p50=15ms, p99=50ms (LOCAL_QUORUM) Instagram observed throughput: RF=1: 150,000 writes/sec RF=3: 145,000 writes/sec (only 3% decrease!) Conclusion: RF=3 performance impact is minimal compared to availability benefits!7What is the difference between RF and number of nodes?Answer: RF and node count are independent concepts!
Key Differences:
Aspect Replication Factor Node Count Definition Copies of each row Total nodes in cluster Scope Per keyspace Cluster-wide Purpose Redundancy/availability Capacity/throughput Typical Value 3 10-1000+ Examples:
# Example 1: 6 nodes, RF=3 - Total nodes: 6 - Copies per row: 3 - Each node stores: 50% of total data (3/6) - Each row appears on: 3 different nodes # Example 2: 300 nodes, RF=3 (Netflix-like) - Total nodes: 300 - Copies per row: 3 - Each node stores: 1% of total data (3/300) - Each row appears on: 3 different nodes (not 300!) # Example 3: 3 nodes, RF=3 - Total nodes: 3 - Copies per row: 3 - Each node stores: 100% of total data (3/3) - Each row appears on: ALL 3 nodes - No capacity benefit, only redundancy! # Example 4: 100 nodes, RF=5 - Total nodes: 100 - Copies per row: 5 - Each node stores: 5% of total data (5/100) - Each row appears on: 5 different nodesRule of Thumb:
Data per node = (Total data × RF) / Total nodes Example with 1TB total data: - 10 nodes, RF=3: 1TB × 3 / 10 = 300GB per node - 30 nodes, RF=3: 1TB × 3 / 30 = 100GB per node - 90 nodes, RF=3: 1TB × 3 / 90 = 33GB per node More nodes = less data per node = better performance!Key Insight: RF controls redundancy/availability, node count controls capacity/performance. Both are important!
AdvertisementResponsive Ad