Section 3: Data Distribution

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.

cqlsh:demo@cassandra

Quick Examples:

cqlsh>
▶ Console ready! Try the quick examples above or write your own commands.
▶ 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.

1️⃣

RF = 1

Copies: 1 (no redundancy)
Fault Tolerance: None
Availability: Low
Use Case: Testing only
Risk: Any node failure = data loss

2️⃣

RF = 2

Copies: 2 (minimal redundancy)
Fault Tolerance: Limited
Availability: Medium
Use Case: Development
Risk: Can't use QUORUM safely

3️⃣

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!

Write with RF=3 (Animated Flow) 👤 Client WRITE user_id=123 Coordinator Node 10.1.0.1 Determines replicas 1. Write R1 10.1.0.1 ✓ Copy 1 Local write (5ms) R2 10.1.0.2 ✓ Copy 2 Network write (12ms) R3 10.1.0.3 ✓ Copy 3 Network write (15ms) 2. Replicate ✅ Write successful! 3 copies stored (RF=3)
# 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

RF Comparison: 1 vs 2 vs 3 (5-Node Cluster) RF = 1 (❌ Risky) N1 N2 ✓ N3 ⚠️ SINGLE POINT OF FAILURE • Data on N2 only (1 copy) • N2 fails → DATA LOST • NO fault tolerance RF = 2 (⚠️ Limited) N1 N2 ✓ N3 ✓ ⚠️ QUORUM PROBLEM • Data on N2 & N3 (2 copies) • QUORUM = 2 (can't tolerate 1 failure) • N2 or N3 fails → QUORUM LOST RF = 3 (✅ Perfect) N1 ✓ N2 ✓ N3 ✓ ✅ FAULT TOLERANT • Data on N1, N2 & N3 (3 copies) • QUORUM = 2 (tolerates 1 failure) • ANY node fails → still QUORUM ✓ Feature Comparison Feature RF = 1 RF = 2 RF = 3 Data Copies 1 2 3 Fault Tolerance 0 nodes 0* nodes 1 node QUORUM Value N/A 2 2 Availability Poor Limited High Storage Cost 1x 2x 3x Write Latency Fastest Medium Slower Production Ready ❌ NO ⚠️ NOT RECOMMENDED ✅ YES Best Use Case Testing only Development Production ⭐ Recommendation: Always use RF=3 for production! ⭐

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!

RF=3 with Different Consistency Levels Node 1 Replica 1 Node 2 Replica 2 Node 3 Replica 3 CL = ONE • Wait for: 1 response • Fault Tolerance: 2 nodes • Latency: Fastest (~5ms) • Use: High availability CL = QUORUM (⭐ Recommended) • Wait for: 2 responses • Fault Tolerance: 1 node • Latency: Balanced (~15ms) • Use: Production standard CL = ALL • Wait for: 3 responses • Fault Tolerance: 0 nodes • Latency: Slowest (~20ms) • Use: Critical data only Key Formulas • QUORUM = (RF / 2) + 1 • Max Failures (QUORUM) = RF - QUORUM • Strong Consistency: R + W > RF (R = Read CL count, W = Write CL count)
# 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 Replication: 3 Datacenters with RF=3 Each 🌎 DC1: US-East (Virginia) RF = 3 (3 copies here) N1 N2 N3 LOCAL_QUORUM = 2 🌎 DC2: US-West (Oregon) RF = 3 (3 copies here) N4 N5 N6 LOCAL_QUORUM = 2 🌍 DC3: EU-West (Ireland) RF = 3 (3 copies here) N7 N8 N9 LOCAL_QUORUM = 2 Total: 9 Copies of Data! ✅ • Each datacenter: RF=3 (3 copies) • Total across cluster: 3 DCs × RF=3 = 9 copies • Can survive: Entire datacenter failure + 1 node in other DC Configuration CREATE KEYSPACE global_data WITH replication = { 'class': 'NetworkTopologyStrategy', 'US-East': 3, # RF=3 in Virginia 'US-West': 3, # RF=3 in Oregon 'EU-West': 3 # RF=3 in Ireland };

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

1️⃣

Always Use RF=3 Minimum

Production Standard: RF=3
Fault Tolerance: 1 node failure
Works With: QUORUM consistency
Used By: Netflix, Instagram, Apple

2️⃣

Use NetworkTopologyStrategy

For: All production clusters
Benefit: DC-aware replication
Allows: Different RF per DC
Future-proof: Easy to add DCs

3️⃣

Pair RF with Correct CL

RF=3: Use CL=QUORUM
Multi-DC: Use LOCAL_QUORUM
R+W > RF: Strong consistency
Balance: Availability vs consistency

4️⃣

Plan for Geographic Distribution

Multiple DCs: Survive region failures
RF per DC: 3 minimum
Racks: Distribute within DC
Result: 99.99%+ availability

5️⃣

Monitor Replication Health

Check: nodetool status regularly
Watch: Inconsistent replicas
Run: nodetool repair weekly
Alert: Node down > 3 hours

6️⃣

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

1
What is Replication Factor and why is it important?
+

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!

2
Why is RF=3 recommended for production?
+

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!

3
How does RF relate to Consistency Level?
+

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!

4
What is NetworkTopologyStrategy and when should you use it?
+

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)
5
How do you check which nodes store replicas for a key?
+

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
6
Can 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 replicas

Best Practice: Plan RF carefully upfront to avoid costly repairs!

7
What 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 restored

Probability:

  • 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!
8
What 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)
6
How 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!
7
What 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 nodes

Rule 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!

Advertisement

Responsive Ad