Section 4: Global Deployments

Multi-Datacenter Deployments

🌍 Master global-scale Cassandra! Learn cross-DC replication, LOCAL_QUORUM, disaster recovery, and latency optimization with real-world examples!

📖 Netflix: 99.99% Global Availability with Multi-DC

Netflix serves 260+ million subscribers globally with Cassandra deployed across 3 AWS regions: US-East, US-West, and EU-West. Configuration: RF=3 per datacenter (9 total copies!). Challenge: When AWS US-East-1 experienced a complete region outage in 2020 affecting 50+ major services, Netflix's viewing experience continued seamlessly. Why? LOCAL_QUORUM consistency level allowed automatic failover to US-West and EU-West. Users in US-East region experienced +80ms latency (querying US-West) but zero service disruption. When US-East recovered 4 hours later, data automatically repaired via hints and anti-entropy. Result: 99.99% uptime maintained, zero data loss, and 140 million concurrent streams unaffected!

💻 Interactive Multi-DC Console

Practice multi-datacenter commands! Configure global replication and consistency levels.

cqlsh:global@cassandra

Quick Examples:

cqlsh>
🌍 Multi-DC console ready!
📡 Try global replication commands
✓ Practice LOCAL_QUORUM and EACH_QUORUM

Learning Tips

  • Global RF: Specify different RF per datacenter
  • LOCAL_QUORUM: Fast local consistency (recommended)
  • EACH_QUORUM: Global strong consistency (slower)
  • Latency: Monitor cross-DC network delays

🎯 Multi-Datacenter Concepts

Multi-datacenter deployment means running Cassandra across multiple geographic locations (datacenters) with automatic data replication between them!

🌍

Geographic Distribution

Purpose: Serve users globally
Benefit: Low latency everywhere
Example: US, EU, Asia regions
Users: Query closest datacenter

🛡️

Disaster Recovery

Purpose: Survive region failures
Benefit: 99.99%+ availability
Example: Entire AWS region down
Result: Service continues uninterrupted

📊

Cost Optimization

Purpose: Balance cost vs availability
Benefit: Different RF per DC
Example: Primary RF=3, Backup RF=2
Savings: Reduce storage costs 30%+

Key Concept

Multi-DC isn't just for "big companies"! Even medium-sized applications benefit from: (1) Lower latency for global users, (2) Protection from cloud provider outages, (3) Compliance with data residency laws (GDPR, etc.). Major companies like Netflix, Uber, Apple all use 3+ datacenters!

🌍 Global Topology Visualization

Global Multi-Datacenter Topology 🇺🇸 US-East RF = 3 (Primary) N1 N2 N3 ✓ 3 copies here 🇺🇸 US-West RF = 3 (Primary) N4 N5 N6 ✓ 3 copies here 🇪🇺 EU-West RF = 2 (Secondary) N7 N8 ✓ 2 copies here Async Replication Cross-Atlantic Global Configuration CREATE KEYSPACE global WITH replication = { 'class': 'NetworkTopologyStrategy', 'US-East': 3, # Primary 'US-West': 3, # Primary 'EU-West': 2 # Secondary }; ✅ Global Benefits ✓ Total Copies: 8 (3+3+2) across 3 continents ✓ Survive: Entire region failure ✓ Latency: <10ms local, ~80-150ms cross-DC
# Complete Multi-DC Configuration Example

# CREATE KEYSPACE with different RF per DC
CREATE KEYSPACE production WITH replication = {
    'class': 'NetworkTopologyStrategy',
    'US-East': 3,      # Primary region - high availability
    'US-West': 3,      # Primary region - high availability
    'EU-West': 2,      # Secondary region - cost optimized
    'Asia-Pacific': 2  # Secondary region - growing market
};

# Total: 10 copies globally!
# US-East: 3 copies (survive 1 failure)
# US-West: 3 copies (survive 1 failure)
# EU-West: 2 copies (survive 0 failures, but US has all data)
# Asia-Pacific: 2 copies (same as EU)

# Check datacenter status
$ nodetool status

Datacenter: US-East
===================
Status=Up/Down
|/ State=Normal/Leaving/Joining/Moving
--  Address      Load       Tokens  Owns   Rack
UN  10.1.0.1     245.5 GB   256     33.3%  us-east-1a
UN  10.1.0.2     238.2 GB   256     33.3%  us-east-1b
UN  10.1.0.3     251.8 GB   256     33.4%  us-east-1c

Datacenter: US-West
===================
UN  10.2.0.1     240.1 GB   256     33.3%  us-west-1a
UN  10.2.0.2     235.9 GB   256     33.3%  us-west-1b
UN  10.2.0.3     248.3 GB   256     33.4%  us-west-1c

Datacenter: EU-West
===================
UN  10.3.0.1     160.5 GB   256     50.0%  eu-west-1a
UN  10.3.0.2     155.2 GB   256     50.0%  eu-west-1b

# All nodes UP and NORMAL!
# Different data sizes due to different RF (EU-West has less)

🔄 Cross-Datacenter Replication

Understanding how data flows between datacenters is critical for performance and consistency!

Cross-DC Write Flow with LOCAL_QUORUM Client US-East User WRITE US-East (Local DC) LOCAL_QUORUM = 2 of 3 N1 ✓ Ack 1 N2 ✓ Ack 2 N3 Async ✅ Success! Returned to client (~5-10ms) Got 2 acks from local DC → LOCAL_QUORUM satisfied Async to US-West US-West +50-80ms later Async to EU-West EU-West +100-150ms later Timeline Breakdown T=0ms: Client writes to coordinator in US-East T=5ms: 2 local nodes ack → LOCAL_QUORUM achieved → SUCCESS returned T=50-150ms: Data async replicates to US-West and EU-West (client doesn't wait)

Why LOCAL_QUORUM is Fast

LOCAL_QUORUM only waits for nodes in the local datacenter. This means: (1) No cross-DC network latency for writes (5-10ms vs 100-150ms), (2) Client gets success immediately, (3) Other DCs receive data asynchronously, (4) Eventually consistent globally!

Consistency Level Comparison

# Consistency Level Behavior in Multi-DC

# LOCAL_QUORUM (Recommended for most writes)
INSERT INTO users VALUES (...) USING CONSISTENCY LOCAL_QUORUM;

Behavior:
- Waits for QUORUM in local DC only
- US-East: Need 2 of 3 local nodes
- Latency: ~5-10ms (local network only)
- Cross-DC: Async (client doesn't wait)
- Use case: 99% of writes

# QUORUM (Global quorum - slower)
INSERT INTO users VALUES (...) USING CONSISTENCY QUORUM;

Behavior:
- Waits for QUORUM of ALL replicas globally
- Total replicas: 8 (3+3+2)
- QUORUM = (8/2)+1 = 5 nodes
- Could need nodes from multiple DCs
- Latency: ~100-150ms (cross-DC network)
- Use case: Rarely needed

# EACH_QUORUM (Strongest consistency - slowest)
INSERT INTO users VALUES (...) USING CONSISTENCY EACH_QUORUM;

Behavior:
- Waits for QUORUM in EACH datacenter
- US-East: Need 2 of 3
- US-West: Need 2 of 3
- EU-West: Need 2 of 2 (both!)
- Latency: ~150-200ms (wait for slowest DC)
- Use case: Critical financial transactions only

# LOCAL_ONE (Fastest - least safe)
INSERT INTO users VALUES (...) USING CONSISTENCY LOCAL_ONE;

Behavior:
- Waits for just 1 node in local DC
- Latency: ~1-2ms
- Risk: Data might not be durable yet
- Use case: High-speed writes, can tolerate loss

⚖️ Multi-DC Consistency Levels

Consistency Levels: Latency vs Consistency Trade-off CL Nodes Needed Latency Use Case LOCAL_ONE 1 node in local DC ~2ms High-speed writes Can tolerate loss LOCAL_QUORUM ⭐ 2 of 3 in local DC ~5-10ms 99% of operations Best balance! QUORUM 5 of 8 (any DCs) ~100ms Cross-DC queries Rarely needed EACH_QUORUM QUORUM in each DC ~150-200ms Financial txns Critical data ALL All 8 replicas ~200ms+ DO NOT USE! Unavailable if 1 node down ⭐ Recommendation: Use LOCAL_QUORUM for 99% of operations!
1️⃣

LOCAL_QUORUM (Default)

Reads: 2 of 3 local nodes
Writes: 2 of 3 local nodes
Latency: 5-10ms
Use: 99% of operations

2️⃣

EACH_QUORUM (Critical)

Reads: QUORUM per DC
Writes: QUORUM per DC
Latency: 150-200ms
Use: Financial transactions only

3️⃣

ALL (Never Use!)

Reads: All 8 nodes
Writes: All 8 nodes
Latency: 200ms+
Problem: 1 node down = unavailable!

⚡ Latency Optimization

Real Numbers: Multi-DC Latency

# Uber's Measured Latencies (Production)

Local DC (LOCAL_QUORUM):
- p50: 5ms
- p95: 12ms
- p99: 25ms
- Success rate: 99.99%

Cross-DC US-East → US-West (QUORUM):
- p50: 85ms
- p95: 120ms
- p99: 180ms
- Success rate: 99.95%

Cross-Atlantic US → EU (QUORUM):
- p50: 140ms
- p95: 200ms
- p99: 300ms
- Success rate: 99.9%

# Key insight: LOCAL_QUORUM is 20-30x faster!
# This is why 99% of operations use LOCAL_QUORUM
⚡

Use Nearest Datacenter

Strategy: Route users to closest DC
Tool: GeoDNS, load balancers
Benefit: Minimize network hops
Result: 5-10ms vs 100-150ms

🎯

LOCAL_QUORUM Everywhere

Reads: LOCAL_QUORUM
Writes: LOCAL_QUORUM
Benefit: No cross-DC waits
Trade-off: Eventually consistent globally

🔧

Tune Connection Pools

Local DC: More connections
Remote DC: Fewer connections
Benefit: Resource optimization
Example: 100 local, 10 remote

🛡️ Disaster Recovery Scenarios

Disaster Recovery: Complete Region Failure ✅ Normal Operation 3 Healthy Datacenters US-East RF=3 US-West RF=3 EU-West RF=2 ⚠️ US-East FAILS Automatic Failover to Remaining DCs US-East ❌ DOWN US-West ✓ SERVING EU-West ✓ SERVING Impact Analysis: Complete Region Failure ✅ What Works ✓ Service continues ✓ Zero data loss ✓ All data available in US-West & EU-West ✓ Automatic failover ✓ No manual intervention ⚠️ Side Effects ⚠️ US-East users routed to US-West ⚠️ Latency +80ms (10ms → 90ms) ⚠️ RF=3 reduced to RF=2 until US-East recovers 🔄 Recovery 🔄 US-East comes back: • Missed writes via hints • Auto-repair starts • Full sync in ~1-2 hrs • Back to normal RF=3

Real Example: AWS US-East-1 Outage (2020)

On November 25, 2020, AWS US-East-1 (Northern Virginia) experienced a major outage affecting Kinesis, Cognito, and other services. Netflix, running Cassandra with RF=3 across US-East, US-West, and EU-West, experienced:

  • Zero downtime: Service continued seamlessly
  • Zero data loss: All data available in US-West & EU-West
  • +80ms latency: US-East users routed to US-West
  • Automatic recovery: When US-East returned after 4 hours, data automatically repaired via hinted handoff

Result: 140 million concurrent streams unaffected, maintaining 99.99% availability!

✅ Multi-DC Best Practices

1️⃣

Minimum 3 Datacenters

Why: Survive entire region failure
Setup: 2 primary (RF=3) + 1 backup (RF=2)
Example: US-East, US-West, EU-West
Benefit: True high availability

2️⃣

Use LOCAL_QUORUM

Reads: LOCAL_QUORUM
Writes: LOCAL_QUORUM
Benefit: Fast (5-10ms)
Exception: EACH_QUORUM for critical only

3️⃣

GeoDNS Routing

Tool: Route53, Cloudflare
Logic: Direct to nearest DC
Benefit: Minimize latency
Failover: Automatic DC switching

4️⃣

Monitor Cross-DC Lag

Metric: Cross-DC replication lag
Threshold: <100ms normally
Alert: >500ms indicates issue
Tool: Prometheus, Grafana

5️⃣

Test Failover Regularly

Frequency: Quarterly
Test: Shut down entire DC
Verify: Service continues
Practice: Chaos engineering

6️⃣

Optimize RF Per Region

Primary: RF=3
Secondary: RF=2
DR Only: RF=1
Savings: 30-40% storage costs

Apple iCloud: Global Multi-DC Setup

# Apple iCloud Cassandra Configuration (Reported)

CREATE KEYSPACE icloud_data WITH replication = {
    'class': 'NetworkTopologyStrategy',
    
    # Primary regions - full redundancy
    'US-East': 3,           # Virginia
    'US-West': 3,           # California
    'EU-West': 3,           # Ireland
    
    # Secondary regions - cost optimized
    'Asia-Pacific': 2,      # Singapore
    'China-North': 3,       # Beijing (compliance)
    'South-America': 2      # São Paulo
};

# Total: 16 copies globally!

# Stats:
- 75,000+ nodes worldwide
- 1+ billion users
- 99.999% availability (5 nines!)
- <10ms latency for local users
- Complete region survival capability

# Consistency:
- Writes: LOCAL_QUORUM (fast, available)
- Reads: LOCAL_QUORUM (fast, consistent locally)
- Critical: EACH_QUORUM (slower but global consistency)

# Result: Seamless global service!

💼 Top 5 Interview Questions

1
What is the difference between LOCAL_QUORUM and QUORUM in a multi-datacenter deployment?
+

Answer:

LOCAL_QUORUM:

  • Scope: Only nodes in the local datacenter
  • Calculation: (Local RF / 2) + 1
  • Example: RF=3 in US-East → LOCAL_QUORUM = 2 nodes in US-East
  • Latency: ~5-10ms (local network only)
  • Cross-DC: Data replicates asynchronously (client doesn't wait)

QUORUM (Global):

  • Scope: All replicas across all datacenters
  • Calculation: (Total RF / 2) + 1
  • Example: RF=3 US-East + RF=3 US-West + RF=2 EU = 8 total → QUORUM = 5
  • Latency: ~100-150ms (may need cross-DC queries)
  • Cross-DC: Must wait for nodes in other datacenters

Detailed Example:

Configuration:
- US-East: RF=3
- US-West: RF=3
- EU-West: RF=2
- Total: 8 replicas

Write from US-East with LOCAL_QUORUM:
1. Client writes to coordinator in US-East
2. Coordinator writes to 3 nodes in US-East
3. Waits for 2 acks (LOCAL_QUORUM = 2 of 3)
4. Returns success to client (~5-10ms)
5. Data asynchronously replicates to US-West and EU-West

Write from US-East with QUORUM:
1. Client writes to coordinator in US-East
2. Coordinator writes to nodes in ALL datacenters
3. Waits for 5 acks total (could be from any DCs)
4. Might need: 2 from US-East + 2 from US-West + 1 from EU
5. Must wait for cross-DC network (~100-150ms)
6. Returns success to client

Why LOCAL_QUORUM is Recommended:

  • Fast: 20-30x lower latency (5ms vs 100ms)
  • Available: Doesn't depend on other DCs being up
  • Consistent: Still provides strong consistency within DC
  • Scalable: Each DC operates independently

When to Use QUORUM:

  • Need global strong consistency immediately (rare)
  • Reading data that might be in different DCs
  • Very specific edge cases (1% of operations)
2
How does Cassandra handle cross-datacenter replication, and is it synchronous or asynchronous?
+

Answer:

Replication Type: Asynchronous

Cross-datacenter replication in Cassandra is asynchronous by default (when using LOCAL_QUORUM). Here's the complete flow:

Write Flow (LOCAL_QUORUM):

T=0ms: Client writes to coordinator in US-East
T=1ms: Coordinator identifies all 8 replica nodes:
       - 3 in US-East (N1, N2, N3)
       - 3 in US-West (N4, N5, N6)
       - 2 in EU-West (N7, N8)

T=2ms: Coordinator writes to LOCAL nodes (US-East):
       - N1: Write starts
       - N2: Write starts
       - N3: Write starts

T=5ms: LOCAL_QUORUM achieved:
       - N1: ✓ Ack received
       - N2: ✓ Ack received
       - N3: Still writing (async)

T=5ms: SUCCESS returned to client!
       ← Client doesn't wait for other DCs!

T=50-80ms: US-West receives data:
       - N4, N5, N6 write asynchronously
       - Client already got success

T=100-150ms: EU-West receives data:
       - N7, N8 write asynchronously
       - Client already got success

Synchronous Option (EACH_QUORUM):

If you need synchronous cross-DC replication, use EACH_QUORUM:

Write with EACH_QUORUM:
T=0ms: Client writes
T=5ms: US-East QUORUM achieved (2 of 3)
T=80ms: US-West QUORUM achieved (2 of 3)
T=150ms: EU-West QUORUM achieved (2 of 2 - both!)
T=150ms: SUCCESS returned to client

Client waits for QUORUM in EVERY datacenter!
Latency: ~150-200ms (slowest DC determines speed)

How Async Replication Works:

  • Coordinator's Job: Coordinator node knows all replica locations via token ring
  • Local First: Writes to local DC nodes immediately
  • Remote Second: Sends writes to remote DC coordinators
  • Remote Coordinators: Each DC has coordinators that handle local writes
  • No Waiting: Client gets success after local writes complete

What if Remote DC is Down?

Scenario: Write with LOCAL_QUORUM, but EU-West is DOWN

T=0ms: Client writes
T=5ms: US-East QUORUM achieved → SUCCESS!
T=50ms: US-West receives data (async)
T=150ms: EU-West unreachable!

Solution: Hinted Handoff
- Coordinator stores "hint" locally
- Hint says: "Send this to EU-West when it's back"
- EU-West comes back online
- Coordinator replays all hints
- EU-West catches up automatically!

Hints stored for: 3 hours by default (configurable)

Consistency Guarantees:

  • LOCAL_QUORUM: Strongly consistent within DC, eventually consistent globally
  • EACH_QUORUM: Strongly consistent globally (but slower)
  • Anti-Entropy: Background repair ensures eventual consistency even without hints

Real-World Usage:

  • Netflix: 99% LOCAL_QUORUM, <1% EACH_QUORUM
  • Uber: All LOCAL_QUORUM except financial transactions
  • Apple: LOCAL_QUORUM for user data, EACH_QUORUM for billing
3
What happens when an entire datacenter goes down in a multi-DC Cassandra deployment?
+

Answer:

Scenario: Complete datacenter failure (power, network, AWS region outage, etc.)

Immediate Effects:

Setup:
- US-East: RF=3 (3 nodes)
- US-West: RF=3 (3 nodes)
- EU-West: RF=2 (2 nodes)
- Total: 8 replicas globally

US-East FAILS (complete region outage):

Immediate (T=0):
✓ US-West: Still has ALL data (3 copies)
✓ EU-West: Still has ALL data (2 copies)
✓ Service continues uninterrupted
✓ Zero data loss

Impact on Users:
⚠️ US-East users: Routed to US-West
   - Latency increases: 5ms → 85ms
   - GeoDNS automatically redirects
   - Service still works, just slower

⚠️ US-West users: No change
   - Still query local DC
   - Latency unchanged: ~5-10ms

⚠️ EU users: No change
   - Still query local DC
   - Latency unchanged: ~5-10ms

What Cassandra Does Automatically:

  • Gossip Protocol: Detects US-East nodes are down within seconds
  • Client Drivers: Automatically stop sending queries to US-East
  • Load Balancer: Routes US-East traffic to US-West
  • Hinted Handoff: US-West and EU-West store hints for US-East
  • No Manual Intervention: Everything automatic!

Consistency Level Behavior:

LOCAL_QUORUM (Recommended):
✓ Writes from US-East users → US-West:
  - US-West has RF=3, needs QUORUM=2
  - Works perfectly!
  - Latency: ~5-10ms in US-West

✓ Reads from US-East users → US-West:
  - US-West has complete data
  - Works perfectly!
  - Latency: ~5-10ms in US-West

✓ US-West and EU users: Unaffected

EACH_QUORUM:
❌ FAILS because can't reach US-East!
   - Needs QUORUM in EACH datacenter
   - US-East unreachable
   - Writes fail until strategy changed

QUORUM (Global):
✓ Still works:
   - Total RF = 8, QUORUM = 5
   - US-West (3) + EU-West (2) = 5
   - Can still satisfy QUORUM
   - But slower than LOCAL_QUORUM

Recovery Process:

When US-East comes back online (T=4 hours):

Phase 1: Nodes Rejoin (T=0-5 minutes)
- US-East nodes start up
- Join gossip network
- Mark themselves as available

Phase 2: Hinted Handoff (T=5-30 minutes)
- US-West replays all hints to US-East
- Catches up on missed writes
- Usually ~10K-100K writes

Phase 3: Anti-Entropy Repair (T=30 mins - 2 hours)
- nodetool repair runs automatically
- Compares data across all DCs
- Fixes any inconsistencies
- Ensures complete synchronization

Phase 4: Back to Normal (T=2-4 hours)
- US-East fully synchronized
- Traffic gradually returns to US-East
- GeoDNS routes US users back
- Latency returns to normal

Real Example: AWS US-East-1 Outage (Nov 2020):

  • Affected: Netflix, Robinhood, Adobe, many others
  • Netflix (Cassandra): Zero downtime
  • Duration: 4 hours
  • Impact: +80ms latency for US-East users, but service continued
  • Recovery: Automatic when US-East returned

Best Practices for DC Failure:

  • Always use LOCAL_QUORUM (EACH_QUORUM fails if any DC down)
  • Minimum 3 DCs: Can lose 1 and still have 2
  • Monitor Gossip: Detect failures immediately
  • Test Failover: Regularly shutdown a DC to verify
  • GeoDNS: Automatic traffic routing
4
How would you design a multi-datacenter Cassandra deployment to minimize latency for global users?
+

Answer: A comprehensive multi-DC design strategy:

1. Datacenter Placement Strategy:

Geographic Coverage:
- Place DCs close to user concentrations
- Example for global app:
  * US-East (Virginia) - 40% of users
  * US-West (California) - 20% of users
  * EU-West (Ireland) - 25% of users
  * Asia-Pacific (Singapore) - 15% of users

Replication Factor per DC:
- High traffic regions: RF=3
- Medium traffic: RF=2
- Low traffic/DR only: RF=1

Example Configuration:
CREATE KEYSPACE global_app WITH replication = {
    'class': 'NetworkTopologyStrategy',
    'US-East': 3,        # 40% users, high traffic
    'US-West': 3,        # 20% users, primary backup
    'EU-West': 3,        # 25% users, high traffic
    'Asia-Pacific': 2    # 15% users, growing
};

2. Client Routing Strategy:

GeoDNS Layer:
1. User connects to app.example.com
2. GeoDNS (Route53/Cloudflare) detects user location
3. Returns IP of nearest datacenter:
   - US users → US-East
   - EU users → EU-West
   - Asia users → Asia-Pacific

Example Route53 Config:
- Record: app.example.com
- Type: Geolocation
- US-East: 10.1.0.1 (health check enabled)
- US-West: 10.2.0.1 (failover for US-East)
- EU-West: 10.3.0.1 (health check enabled)
- Asia-Pacific: 10.4.0.1 (health check enabled)

Result: User always connects to nearest DC (5-10ms latency)

3. Consistency Level Strategy:

Default (99% of operations):
- Reads: LOCAL_QUORUM
- Writes: LOCAL_QUORUM
- Latency: 5-10ms
- Consistency: Strong within DC, eventual across DCs

Critical Operations (1%):
- Financial transactions: EACH_QUORUM
- Billing updates: EACH_QUORUM
- Latency: 150-200ms (acceptable for critical ops)

Application Code:
// Default queries
session.execute(query, ConsistencyLevel.LOCAL_QUORUM);

// Critical financial transaction
if (isCriticalTransaction) {
    session.execute(query, ConsistencyLevel.EACH_QUORUM);
}

4. Connection Pool Optimization:

DataStax Java Driver Configuration:

// Application running in US-East
Cluster cluster = Cluster.builder()
    .addContactPoints("us-east-node1", "us-east-node2")
    .withLoadBalancingPolicy(
        DCAwareRoundRobinPolicy.builder()
            .withLocalDc("US-East")              // Prefer local DC
            .withUsedHostsPerRemoteDc(2)         // Only 2 remote connections
            .build()
    )
    .withPoolingOptions(
        new PoolingOptions()
            .setConnectionsPerHost(HostDistance.LOCAL, 4, 10)   // Local: 4-10 connections
            .setConnectionsPerHost(HostDistance.REMOTE, 1, 2)   // Remote: 1-2 connections
    )
    .build();

Result:
- Queries go to local DC 99.9% of time
- Remote DCs only for failover
- Minimizes cross-DC traffic

5. Read/Write Path Optimization:

  • Write-heavy workloads: LOCAL_QUORUM writes, async replication to other DCs
  • Read-heavy workloads: LOCAL_QUORUM reads from nearest DC
  • Mixed workloads: Partition data by geography when possible

6. Caching Strategy:

Multi-Layer Caching:
1. Application Cache (Redis/Memcached) - per DC
   - Hot data: 5-minute TTL
   - Latency: <1ms

2. Cassandra Row Cache - per DC
   - Frequently accessed rows
   - Latency: 1-2ms

3. Cassandra Query - LOCAL_QUORUM
   - Cache miss
   - Latency: 5-10ms

Result: 90%+ requests served from cache (<1ms)

7. Monitoring & Alerting:

Key Metrics to Monitor:
- Per-DC latency (p50, p95, p99)
  * Alert if p99 > 50ms for LOCAL_QUORUM
- Cross-DC replication lag
  * Alert if lag > 500ms
- Per-DC error rate
  * Alert if errors > 0.1%
- Hint pressure
  * Alert if hints queue > 10K

Tools:
- Prometheus + Grafana
- DataStax OpsCenter
- Custom dashboards per DC

Expected Latency Results:

Operation p50 p95 p99
Cache hit 0.5ms 1ms 2ms
LOCAL_QUORUM read 5ms 12ms 25ms
LOCAL_QUORUM write 6ms 15ms 30ms
5
What are the trade-offs between using different RF values per datacenter?
+

Answer: Understanding RF trade-offs per DC is crucial for cost and availability balance:

Scenario Comparison:

Scenario Configuration Total Copies Storage Cost
Uniform High US-East:3, US-West:3, EU:3 9 9x data
Mixed (Recommended) US-East:3, US-West:3, EU:2 8 8x data (-11%)
Cost Optimized US-East:3, US-West:2, EU:1 6 6x data (-33%)

Detailed Analysis:

1. Primary DC (RF=3) - High Traffic Region:

Benefits:
✓ Can lose 1 node and maintain QUORUM (2 of 3)
✓ Strong consistency with LOCAL_QUORUM
✓ 99.99% availability
✓ Best for user-facing DCs

Trade-offs:
✗ Higher storage cost (3x data in this DC)
✗ More nodes to manage
✗ Higher operational complexity

Use when:
- High user traffic (US-East, EU-West)
- Primary serving regions
- Critical to business

2. Secondary DC (RF=2) - Medium Traffic:

Benefits:
✓ Moderate storage cost (2x data in this DC)
✓ Can still serve with LOCAL_QUORUM (2 of 2)
✓ Good for disaster recovery
✓ Cost savings vs RF=3

Trade-offs:
✗ CANNOT lose any nodes without losing QUORUM
✗ LOCAL_QUORUM requires both nodes (less fault tolerant)
✗ If 1 node down, must use CL=ONE (weaker consistency)
✗ Lower availability than RF=3

Warning with RF=2:
If 1 of 2 nodes down:
- LOCAL_QUORUM needs 2 of 2 → FAILS!
- Must downgrade to CL=ONE
- Reduced consistency guarantee

Use when:
- Secondary regions with less traffic
- Disaster recovery backup
- Cost optimization needed
- Can tolerate reduced availability

3. DR-Only DC (RF=1) - Rarely Used:

Benefits:
✓ Lowest storage cost (1x data in this DC)
✓ Complete data copy for disaster recovery
✓ Good for compliance (data in specific region)

Trade-offs:
✗ NO redundancy within this DC
✗ Cannot use LOCAL_QUORUM (need 1, have 1)
✗ Must use CL=ONE (weakest consistency)
✗ Node failure = data unavailable in this DC
✗ Should NOT serve user traffic

Use when:
- Disaster recovery only (not serving traffic)
- Compliance requirement (must have data in region)
- Extreme cost optimization
- Other DCs have full RF=3

DO NOT use RF=1 if:
- Serving user traffic in this DC
- Need any fault tolerance
- PRIMARY serving region

Real-World Example: Netflix Configuration:

CREATE KEYSPACE viewing_history WITH replication = {
    'class': 'NetworkTopologyStrategy',
    'US-East': 3,      # Primary - 60% of users
    'US-West': 3,      # Primary - backup and 20% users
    'EU-West': 3       # Primary - 20% of users
};
# Total: 9 copies (expensive but maximum availability)

CREATE KEYSPACE recommendation_cache WITH replication = {
    'class': 'NetworkTopologyStrategy',
    'US-East': 3,      # Primary
    'US-West': 2,      # Secondary (can regenerate)
    'EU-West': 2       # Secondary
};
# Total: 7 copies (22% cost savings, acceptable for cache)

CREATE KEYSPACE audit_logs WITH replication = {
    'class': 'NetworkTopologyStrategy',
    'US-East': 3,      # Primary
    'US-West': 2,      # Backup
    'EU-West': 1       # Compliance only
};
# Total: 6 copies (33% cost savings, rarely queried)

Cost Impact Example:

Scenario: 1TB of raw data

All RF=3 (Uniform):
- US-East: 3TB
- US-West: 3TB
- EU-West: 3TB
- Total: 9TB storage
- Monthly cost (AWS): ~$900 (9TB × $100/TB)

Mixed RF (Recommended):
- US-East: 3TB (RF=3)
- US-West: 3TB (RF=3)
- EU-West: 2TB (RF=2)
- Total: 8TB storage
- Monthly cost: ~$800 (8TB × $100/TB)
- Savings: $100/month, $1,200/year per TB

Cost Optimized:
- US-East: 3TB (RF=3)
- US-West: 2TB (RF=2)
- EU-West: 1TB (RF=1)
- Total: 6TB storage
- Monthly cost: ~$600 (6TB × $100/TB)
- Savings: $300/month, $3,600/year per TB

But: Reduced availability in US-West and EU-West!

Decision Matrix:

  • Use RF=3: Primary user-facing DCs, critical data, high traffic
  • Use RF=2: Secondary DCs, backup regions, cost-conscious deployments
  • Use RF=1: DR only, compliance requirements, NEVER for serving traffic

Key Principle: Always have at least 2 DCs with RF≥3 to survive complete DC failure!

Advertisement

Responsive Ad