Multi-Datacenter Consistency

LOCAL - Multi-Datacenter Mastery

๐ŸŒ Master LOCAL_QUORUM & LOCAL_ONE! Eliminate cross-datacenter latency while maintaining strong consistency!

๐Ÿ“– Uber: When LOCAL_QUORUM Saves Millions

Uber operates in 10,000+ cities across 70+ countries with multiple datacenters globally. For ride requests, Uber needs strong consistency (can't have duplicate rides!) but also low latency (drivers/riders waiting). Challenge: Using regular QUORUM across datacenters = 200-500ms cross-continent latency! A ride request in San Francisco shouldn't wait for Tokyo datacenter! Uber's solution: LOCAL_QUORUM - get strong consistency within the local datacenter only (~5-10ms). Result: Sub-15ms latency, riders matched instantly, drivers respond quickly. Trade-off: If San Francisco DC fails, Tokyo DC might have slightly stale data temporarily. Impact: Acceptable! Cross-DC replication happens async (within 100-200ms). Benefit: 20-50x faster writes compared to cross-DC QUORUM, millions saved in infrastructure, better user experience globally. Uber engineers: "For multi-DC deployments, LOCAL_QUORUM is essential - you get strong consistency locally without cross-continent latency!"

๐Ÿ’ป Interactive LOCAL Simulator

See how LOCAL consistency levels work across multiple datacenters in real-time!

cqlsh@multi-dc-cluster

๐ŸŒ Multi-DC Scenarios - Click to Explore:

๐ŸŒ Multi-DC LOCAL Simulator Ready!
๐ŸŽฏ Click scenarios above to see LOCAL in action
โœ“ Learn LOCAL_QUORUM vs LOCAL_ONE across datacenters

๐ŸŒ Why LOCAL Matters

CRITICAL: If you're running Cassandra across multiple datacenters, LOCAL consistency levels are essential! Without LOCAL, every write waits for cross-datacenter confirmation = 200-500ms latency. With LOCAL_QUORUM, you get strong consistency locally (~10ms) while async replication handles cross-DC sync. This is how Netflix, Uber, and Apple run global Cassandra deployments!

๐ŸŒ What are LOCAL Consistency Levels?

Imagine you run 3 coffee shops - one in New York, one in London, and one in Tokyo. You want to track inventory. With LOCAL_QUORUM, when the New York shop updates inventory, you only wait for confirmation from 2 shops in New York's region, not all 3 globally. You don't wait for Tokyo's response (which would take 200ms+). Tokyo gets the update asynchronously in the background.

๐ŸŽฏ Real-World Analogy: Regional Managers

Your company has 3 regional offices (US, Europe, Asia). Each region has 3 managers. You need to make a decision.

With Regular QUORUM (Cross-Region):
You call a meeting with 2 managers from ANY region globally. Problem: If you're in US and call Europe + Asia managers, you wait for international calls = 300ms delay! Slow!

With LOCAL_QUORUM (Same Region):
You only meet with 2 managers from YOUR region (US). Decision made in 10ms! Fast! Europe and Asia managers get notified later (async). They catch up within 100-200ms.

Perfect for: Decisions that need regional consistency but not immediate global sync
Trade-off: If US region fails, Europe/Asia might have slightly stale data temporarily

This is exactly how LOCAL_QUORUM works in multi-datacenter Cassandra!

๐ŸŒ

LOCAL_QUORUM

Behavior: QUORUM within local datacenter only
Latency: ~5-10ms (no cross-DC wait)
Consistency: Strong locally, eventual across DCs
What happens: Wait for 2 of 3 local replicas
Use when: Multi-DC + need strong consistency locally

โšก

LOCAL_ONE

Behavior: ONE within local datacenter only
Latency: ~3-5ms (fastest local)
Consistency: Eventual locally AND across DCs
What happens: Wait for 1 local replica only
Use when: Multi-DC + speed priority + eventual OK

๐Ÿ”„

Cross-DC Replication

Mechanism: Async background sync
Timeline: Usually 100-200ms across continents
Guarantee: All DCs eventually consistent
Benefit: No blocking on cross-DC latency
Perfect for: Global applications

๐Ÿ’ก Key Insight: LOCAL = Speed + Multi-DC

LOCAL consistency levels solve the "multi-datacenter latency problem". Without LOCAL:

  • QUORUM in multi-DC: Wait for 2 replicas from ANY DC = could be US + Tokyo = 200ms+
  • ALL in multi-DC: Wait for ALL replicas globally = disaster (500ms+ or timeout)

With LOCAL:

  • LOCAL_QUORUM: Wait for 2 replicas in YOUR DC = ~10ms (fast!)
  • LOCAL_ONE: Wait for 1 replica in YOUR DC = ~5ms (fastest!)
  • Cross-DC sync: Happens async, no blocking

This is why every major company with multi-DC Cassandra uses LOCAL consistency levels!

๐Ÿ”ง How LOCAL_QUORUM Works Across Datacenters

Write with LOCAL_QUORUM: Only Wait for Local DC (RF=3 per DC) ๐Ÿ‡บ๐Ÿ‡ธ Client (US) Writes with LOCAL_QUORUM ๐Ÿข US Datacenter (Local) US Replica 1 โšก โœ“ ACK T = 5ms US Replica 2 โšก โœ“ ACK T = 8ms ๐Ÿข Europe DC (Remote - Async) EU Replica (3 nodes) โณ Async T = 150ms (background) SUCCESS at ~8ms! Async replication โฑ๏ธ Complete Timeline Comparison โœ… LOCAL_QUORUM (Smart!) T=0ms: Client sends write T=5ms: US Replica 1 ACKs T=8ms: US Replica 2 ACKs โ†’ SUCCESS at 8ms! โšก T=150ms: EU updated (async) โŒ Regular QUORUM (Slow!) T=0ms: Client sends write T=5ms: US Replica 1 ACKs T=8ms: US Replica 2 ACKs T=150ms: Wait for EU... โ†’ SUCCESS at 150ms ๐ŸŒ ๐ŸŽฏ Key Insight: LOCAL_QUORUM is 18x faster! Strong consistency locally + No cross-DC wait = Best of both worlds for multi-DC deployments

โšก Performance Magic: 18x Faster!

Notice the difference? LOCAL_QUORUM took 8ms while regular QUORUM took 150ms - that's 18x faster! The key insight: LOCAL_QUORUM gets strong consistency by waiting for 2 local replicas, but doesn't wait for cross-datacenter acknowledgments. Europe datacenter gets the update asynchronously (within 150ms) without blocking the client. This is how Uber, Netflix, and Apple achieve sub-15ms writes globally!

๐ŸŒ The Cross-Datacenter Latency Problem

Cross-datacenter network latency is physics - you can't avoid it! Light/data travels at ~200,000 km/second through fiber. Here's what that means:

๐ŸŒ

US East โ†’ US West

Distance: ~4,000 km
Latency: 60-80ms round-trip
Impact: QUORUM across coasts = 60-80ms writes
With LOCAL: 5-10ms (12x faster!)
Use case: E-commerce across US regions

๐ŸŒŽ

US โ†’ Europe

Distance: ~6,000 km
Latency: 120-150ms round-trip
Impact: QUORUM = 150ms writes (unacceptable!)
With LOCAL: 5-10ms (20x faster!)
Use case: Global SaaS applications

๐ŸŒ

US โ†’ Asia

Distance: ~10,000+ km
Latency: 200-300ms round-trip
Impact: QUORUM = 300ms writes (disaster!)
With LOCAL: 5-10ms (30x faster!)
Use case: Global gaming, ride-sharing

๐Ÿ’ฅ Real Production Horror Story

A major e-commerce company launched with Cassandra across 3 datacenters (US, EU, Asia). They used regular QUORUM for all writes (didn't know about LOCAL!). Result: Every product add-to-cart took 200-300ms because QUORUM waited for Asia datacenter ACK! Users complained about slow checkout. Cart abandonment increased 40%. Engineers discovered LOCAL_QUORUM, switched to it, latency dropped to 8ms. Cart abandonment back to normal. Lesson: If you have multiple DCs, you MUST use LOCAL consistency levels!

โœ… LOCAL_QUORUM: The Multi-DC Standard

LOCAL_QUORUM is the recommended default for multi-datacenter deployments. It gives you strong consistency within your local DC plus the performance benefits of avoiding cross-DC latency.

โœ…

When to Use LOCAL_QUORUM

Perfect for: 99% of multi-DC production workloads!
User data: Profiles, accounts, settings
Business logic: Orders, inventory, reservations
Benefit: Strong consistency locally, ~8-10ms latency
Trade-off: Eventual consistency across DCs
Real use: Uber rides, Netflix accounts, Apple services

๐Ÿ”’

Consistency Guarantee

Local DC: Strong consistency (QUORUM of local replicas)
Formula: R_local + W_local > RF_local
Example: RF=3 per DC โ†’ 2 + 2 = 4 > 3 โœ“
Guarantee: Reads see latest local writes
Cross-DC: Eventually consistent (~100-200ms)

โšก

Performance Characteristics

Write latency: 8-12ms (same-DC QUORUM)
Read latency: 8-12ms (no cross-DC wait)
Throughput: 30,000+ writes/sec per node
vs QUORUM: 15-20x faster in multi-DC
CPU usage: 65% vs 95% with cross-DC QUORUM

โš ๏ธ

Important Caveats

DC failure: Writes to failed DC will fail
Mitigation: Client failover to healthy DCs
Cross-DC lag: 100-300ms replication delay
Solution: Use EACH_QUORUM for critical cross-DC
Network partitions: Each DC operates independently

๐Ÿ“Š

Configuration Example

-- Write with LOCAL_QUORUM
INSERT INTO users (id, name, email)
VALUES (123, 'Alice', 'alice@example.com')
USING CONSISTENCY LOCAL_QUORUM;

-- Read with LOCAL_QUORUM
SELECT * FROM users WHERE id=123
USING CONSISTENCY LOCAL_QUORUM;
๐Ÿข

Multi-DC Setup

Typical: 3 DCs (US, EU, Asia)
RF per DC: 3 (total 9 replicas globally)
Strategy: NetworkTopologyStrategy
Config: {'US':3, 'EU':3, 'Asia':3}
Client routing: DC-aware load balancer

โšก LOCAL_ONE: Speed in Multi-DC

LOCAL_ONE is like ONE, but ensures you read/write from your local datacenter only. Perfect for high-volume, non-critical data in multi-DC setups.

โšก

LOCAL_ONE vs ONE

LOCAL_ONE: 1 replica from local DC only
ONE: 1 replica from ANY DC (could be remote!)
Problem with ONE: Might read from Asia DC = 200ms
LOCAL_ONE benefit: Always local = ~3-5ms
Use case: Analytics, logs, metrics in multi-DC

๐Ÿ“Š

When to Use LOCAL_ONE

Perfect for: High-volume non-critical in multi-DC
Analytics: Pageviews, clicks, impressions
Logs: Application logs, event streams
Caches: Session cache, API response cache
Benefit: 3-5ms writes, 50k+ writes/sec
Trade-off: Eventual consistency locally AND globally

โŒ

When NOT to Use LOCAL_ONE

Never for: User-facing data in multi-DC
User profiles: Use LOCAL_QUORUM
Orders: Use LOCAL_QUORUM
Authentication: Use LOCAL_QUORUM
Why dangerous: Stale reads WITHIN same DC!
Rule: If ONE is dangerous, LOCAL_ONE is too

โšก Performance Benefits: Numbers Don't Lie

๐Ÿ“Š Real Production Benchmarks

Setup: 3 datacenters (US East, Europe, Asia), RF=3 per DC, 1 million writes from US client

LOCAL_ONE (Fastest):
โ€ข Latency: P50=3ms, P99=6ms
โ€ข Throughput: 50,000 writes/sec
โ€ข Total time: 20 seconds
โ€ข CPU: 40% average

LOCAL_QUORUM (Recommended):
โ€ข Latency: P50=8ms, P99=15ms
โ€ข Throughput: 30,000 writes/sec
โ€ข Total time: 33 seconds
โ€ข CPU: 65% average

QUORUM (Cross-DC - Slow!):
โ€ข Latency: P50=150ms, P99=300ms
โ€ข Throughput: 2,000 writes/sec
โ€ข Total time: 500 seconds (8+ minutes!)
โ€ข CPU: 95% average

ALL (Cross-DC - Disaster!):
โ€ข Latency: P50=280ms, P99=500ms
โ€ข Throughput: 1,000 writes/sec
โ€ข Total time: 1000 seconds (16+ minutes!)
โ€ข CPU: 98% average

Performance Gains:
โ€ข LOCAL_QUORUM vs QUORUM: 18x faster, 15x more throughput
โ€ข LOCAL_QUORUM vs ALL: 35x faster, 30x more throughput
โ€ข LOCAL_ONE vs QUORUM: 50x faster, 25x more throughput

๐Ÿข Real Companies Using LOCAL Consistency

๐Ÿš— Uber: Ride Matching at Global Scale

Scale: 10,000+ cities, 70+ countries, multiple DCs globally
Use Case: Ride requests and driver matching
Why LOCAL_QUORUM: Need strong consistency for matching (can't duplicate rides!) but can't wait 200ms for cross-DC confirmation
Configuration: Write rides with LOCAL_QUORUM in nearest DC
Latency Impact: Sub-15ms ride matching vs 200-300ms with cross-DC QUORUM
Result: Instant rider-driver matching, drivers respond quickly, better UX globally
Performance Gain: 20-30x faster than cross-DC QUORUM
Trade-off: If US DC fails, US rides temporarily unavailable (acceptable - rare event)
Uber Engineering: "LOCAL_QUORUM is essential for multi-DC deployments. Without it, ride matching would be too slow globally."

๐Ÿ“บ Netflix: User Accounts & Subscriptions

Scale: 260M+ subscribers across 190+ countries, DCs worldwide
Use Case: User account data, subscription status, payment methods
Why LOCAL_QUORUM: User account updates need strong consistency but <10ms latency for good UX
Hybrid Approach:
โ€ข LOCAL_QUORUM: User profiles, account settings, preferences
โ€ข LOCAL_ONE: Viewing history, watch progress (billions/day)
โ€ข QUORUM: Billing, payment processing (critical cross-DC)
Impact: Sub-10ms account updates globally, billions saved in infrastructure
Netflix Tech Blog: "LOCAL_QUORUM lets us maintain strong consistency where it matters while avoiding the latency penalty of cross-DC coordination."

๐ŸŽ Apple: iCloud Services

Scale: 2B+ active devices, datacenters globally
Use Case: iCloud data sync, device backups, user settings
Why LOCAL_QUORUM: Users expect instant sync within their region, eventual global sync acceptable
Configuration:
โ€ข 5 datacenters globally (US West, US East, EU, Asia Pacific, Middle East)
โ€ข RF=3 per DC, LOCAL_QUORUM for all user data writes
โ€ข Cross-DC async replication (~100-200ms)
User Experience: Settings change on iPhone in California โ†’ instant confirmation โ†’ synced to Europe within 150ms
Performance: 5-10ms local writes vs 200-400ms if using cross-DC QUORUM
Apple's Choice: LOCAL_QUORUM + eventual cross-DC = perfect balance for global consumer services

โœ… Best Practices for LOCAL Consistency Levels

๐ŸŽฏ

1. Always Use LOCAL in Multi-DC

Rule: If you have >1 DC, use LOCAL_*
Default: LOCAL_QUORUM for production
Never: Use QUORUM/ALL in multi-DC
Exception: EACH_QUORUM for critical cross-DC
Testing: Verify with RF=3 per DC in staging

๐Ÿ”„

2. Hybrid Consistency Strategy

Critical data: LOCAL_QUORUM (users, orders)
High-volume: LOCAL_ONE (analytics, logs)
Cross-DC critical: EACH_QUORUM (billing)
Pattern: Different CL per table/query
Benefit: Optimize each use case

๐Ÿข

3. DC-Aware Client Routing

Setup: Route clients to nearest DC
Example: US users โ†’ US DC, EU users โ†’ EU DC
Benefit: LOCAL_QUORUM = truly local
Failover: Route to healthy DC if local fails
Tool: DC-aware load balancer

๐Ÿ“Š

4. Monitor Cross-DC Lag

Metric: Replication lag between DCs
Target: <200ms for most workloads
Alert: If lag >500ms consistently
Tool: Cassandra's nodetool netstats
Fix: Check network, increase bandwidth

๐Ÿ’พ

5. Plan for DC Failures

Expect: DCs will fail (rare but happens)
Impact with LOCAL: Writes to failed DC fail
Mitigation: Client auto-failover to healthy DC
Recovery: Hinted handoff when DC returns
Test: Run DR drills regularly

๐Ÿงช

6. Test Cross-DC Scenarios

Test 1: Write in US, read in EU (eventual lag)
Test 2: US DC fails, EU clients work
Test 3: Network partition between DCs
Test 4: Cross-DC repair after outage
Tool: Chaos engineering (Chaos Monkey)

โš ๏ธ Common Mistakes to Avoid

  • โŒ Using QUORUM in multi-DC: Results in 150-300ms cross-DC latency
  • โŒ Forgetting LOCAL_ prefix: ONE might read from remote DC!
  • โŒ Not testing DC failures: Surprises in production when DC fails
  • โŒ Same CL for everything: Use hybrid strategy per use case
  • โŒ Ignoring cross-DC lag: Monitor and alert on replication delays
  • โŒ No client failover: Single DC failure = total outage

๐Ÿ’ผ Interview Questions & Answers

1
What is LOCAL_QUORUM and why is it essential for multi-datacenter deployments?

Complete Answer:

LOCAL_QUORUM is a consistency level that requires QUORUM (majority) of replicas to acknowledge within the local datacenter only, not globally. It's the recommended default for multi-DC Cassandra deployments.

How it works:

  • Write: Coordinator sends write to all replicas globally but only waits for QUORUM of local DC replicas to ACK before returning success
  • Example: RF=3 per DC โ†’ Wait for 2 ACKs from local DC only (~8ms), remote DCs updated asynchronously (~150ms)
  • Read: Contact QUORUM of replicas in local DC only

Why essential for multi-DC:

  • Cross-DC latency problem: USโ†’Europe = 150ms, USโ†’Asia = 250-300ms (physics - can't avoid!)
  • With regular QUORUM: Need 2 ACKs from ANY DCs = might wait for Asia = 300ms (disaster!)
  • With LOCAL_QUORUM: Need 2 ACKs from local DC only = ~8ms (18-35x faster!)
  • Strong consistency locally: R_local + W_local > RF_local = 2 + 2 = 4 > 3 (guaranteed fresh reads locally)
  • Eventual across DCs: Cross-DC replication happens async (100-200ms) - acceptable for most use cases

Real production impact:

  • Uber: Ride matching ~10ms with LOCAL_QUORUM vs 200-300ms with cross-DC QUORUM
  • Netflix: Account updates ~8ms locally, billions saved in infrastructure
  • Apple: iCloud sync instant within region, eventual globally

Key insight: In multi-DC deployments, LOCAL_QUORUM gives you the best of both worlds: strong consistency where you need it (local DC) and performance that doesn't suffer from cross-continent network latency. This is why virtually every major company with global Cassandra deployments uses LOCAL_QUORUM as their default consistency level!

2
Explain the difference between ONE and LOCAL_ONE. Why does LOCAL_ONE exist?

Complete Answer:

Both ONE and LOCAL_ONE wait for only 1 replica to acknowledge, but they differ crucially in which replica that can be.

Key Difference:

  • ONE: Waits for 1 replica from ANY datacenter globally (could be local or remote!)
  • LOCAL_ONE: Waits for 1 replica from local datacenter only (guaranteed local)

The Problem with ONE in Multi-DC:

  • Scenario: You're a US client reading with CL=ONE from a cluster with DCs in US, Europe, and Asia
  • Problem: Cassandra might randomly pick a replica in Asia datacenter!
  • Result: Read latency = 250-300ms (cross-Pacific network) even though there's a replica right next to you in US
  • Impact: Unpredictable latency - sometimes 3ms, sometimes 300ms!

Why LOCAL_ONE exists:

  • Guarantee: Always reads from local DC = predictable low latency (~3-5ms)
  • Use case: High-volume non-critical data in multi-DC (analytics, logs, metrics)
  • Benefit: Fastest possible local reads, no cross-DC network hops
  • Trade-off: Eventual consistency both locally AND across DCs

Detailed Example:

Setup: 3 DCs (US, EU, Asia), RF=3 per DC, US client

  • With ONE:
    • Cassandra picks random replica globally
    • 33% chance: US replica = 3ms โœ“
    • 33% chance: EU replica = 150ms โŒ
    • 33% chance: Asia replica = 300ms โŒ
    • P50 latency: ~150ms (unpredictable!)
  • With LOCAL_ONE:
    • Always picks US replica (local DC)
    • 100% chance: US replica = 3-5ms โœ“
    • P50 latency: 3ms (predictable!)

When to use each:

  • Use LOCAL_ONE when: Multi-DC deployment + high-volume non-critical data (analytics, logs)
  • Use ONE when: Single DC deployment (LOCAL_ONE = ONE in single DC)
  • Never use ONE in multi-DC: Unpredictable cross-DC latency!

Production recommendation: If you have multiple datacenters, always use LOCAL_ variants (LOCAL_ONE, LOCAL_QUORUM) instead of ONE or QUORUM. This ensures predictable latency and avoids unexpected cross-datacenter network hops!

3
What happens if a datacenter fails when using LOCAL_QUORUM? How do you handle this?

Complete Answer:

When a datacenter fails with LOCAL_QUORUM, writes/reads to that specific datacenter will fail, but other datacenters continue operating normally. This is an important trade-off to understand.

What Happens During DC Failure:

  • Failed DC: Clients trying to write/read with LOCAL_QUORUM to the failed DC will get errors (can't get QUORUM from a down DC)
  • Other DCs: Continue working normally - clients routed to healthy DCs operate fine
  • Example: US DC fails โ†’ US clients get errors, EU/Asia clients work normally

Contrast with regular QUORUM:

  • With QUORUM: If US DC fails, can still get QUORUM from EU + Asia DCs = writes succeed
  • Better availability during DC failures
  • BUT: Normal operation = 150-300ms latency (unacceptable!)
  • Trade-off: 99.9% of time speed vs 0.1% of time availability

How to Handle DC Failures with LOCAL_QUORUM:

1. Client-Side Failover (Most Common):

  • Setup: DC-aware load balancer (e.g., AWS Route 53, Consul)
  • Normal operation: US clients โ†’ US DC, EU clients โ†’ EU DC
  • On US DC failure: Load balancer detects failure, routes US clients to EU DC
  • Result: US clients now write to EU DC with LOCAL_QUORUM = works! (slightly higher latency ~150ms but better than complete outage)
  • Recovery: When US DC comes back, clients gradually fail back

2. Application-Level Retry:

// Example retry logic
try {
    session.execute(query, LOCAL_QUORUM); // Try local DC
} catch (NoHostAvailableException e) {
    // Local DC down, retry with QUORUM to fail over to other DCs
    session.execute(query, QUORUM);
}

3. Use EACH_QUORUM for Critical Writes:

  • For critical data: Use EACH_QUORUM instead of LOCAL_QUORUM
  • EACH_QUORUM: Requires QUORUM in EACH datacenter
  • Benefit: If one DC fails, other DCs still have quorum of data
  • Cost: Slower (cross-DC latency) but better availability
  • Use case: Billing, payment processing, critical business data

4. Hinted Handoff for Recovery:

  • During outage: Cassandra stores "hints" for failed DC
  • When DC recovers: Hints are replayed to catch up missed writes
  • Timeline: Usually catches up in minutes to hours
  • Limitation: Hints only stored for 3 hours by default (configurable)

5. Monitor DC Health:

  • Metrics: Monitor DC availability, latency, error rates
  • Alerts: Immediate notification on DC failures
  • Tools: DataStax OpsCenter, Prometheus + Grafana
  • SLOs: Set DC availability SLOs (e.g., 99.95%)

Production Decision Framework:

  • Most companies choose: LOCAL_QUORUM + client failover
  • Reasoning: 99.9% of time = 10ms latency, 0.1% of time (DC failure) = 150ms + failover
  • Alternative: Always using QUORUM = 100% of time = 150ms (unacceptable for user experience)
  • Critical data: Use EACH_QUORUM for must-never-lose data (billing, compliance)

Real-world example: When AWS US-EAST-1 had a major outage, companies with LOCAL_QUORUM + proper failover automatically routed traffic to US-WEST and EU datacenters. Users experienced slightly higher latency but service stayed up. Companies without failover had complete outages in affected regions.

4
How does cross-datacenter replication work with LOCAL_QUORUM? Explain the timeline.

Complete Answer:

With LOCAL_QUORUM, writes are acknowledged locally fast, then replicated to remote datacenters asynchronously in the background. Understanding this timeline is crucial for multi-DC deployments.

Complete Write Timeline (US client, 3 DCs: US, EU, Asia):

T=0ms:

  • Client in US sends write with CL=LOCAL_QUORUM
  • Coordinator node in US receives write
  • Coordinator determines which replicas need the data (3 in US, 3 in EU, 3 in Asia = 9 total)

T=0-8ms (Local DC writes):

  • Coordinator sends write to all 3 US replicas in parallel
  • T=5ms: US Replica 1 ACKs โœ“
  • T=8ms: US Replica 2 ACKs โœ“
  • LOCAL_QUORUM satisfied! (2 of 3 local replicas)
  • Coordinator returns SUCCESS to client
  • Client latency: 8ms (user sees instant confirmation)

T=8-150ms (Cross-DC replication - Async):

  • Coordinator already sent write to EU/Asia coordinators at T=0
  • EU coordinator writes to 3 EU replicas (happens in parallel with US writes)
  • Asia coordinator writes to 3 Asia replicas (happens in parallel)
  • T=150ms: EU replicas have data โœ“ (USโ†’EU network = ~120-150ms)
  • T=280ms: Asia replicas have data โœ“ (USโ†’Asia network = ~250-300ms)
  • Client doesn't wait for EU/Asia - already got success at T=8ms!

T=280ms+: Global Consistency Achieved:

  • All 9 replicas globally now have the data
  • System is globally consistent (eventual consistency complete)
  • Total replication time: ~280ms
  • But client latency was only 8ms!

What Happens During Async Replication:

  • Coordinator role: US coordinator sends write to EU/Asia coordinators immediately (T=0)
  • Remote coordinators: Each DC's coordinator writes to its local replicas independently
  • No blocking: US client doesn't wait for EU/Asia coordinators to finish
  • Parallel writes: All DCs receive and process write simultaneously
  • Network latency: Cross-DC network adds 150-300ms but doesn't block client

Reading During Replication Window:

Scenario: US client writes at T=0, EU client reads at T=50ms

  • Problem: EU DC hasn't received data yet (arrives at T=150ms)
  • EU read with LOCAL_QUORUM: Returns old data (stale!) โŒ
  • This is expected! Eventual consistency across DCs
  • By T=150ms: EU reads return new data โœ“

Replication Mechanics:

  • Protocol: Cassandra's inter-DC streaming
  • Optimization: Batching (multiple writes bundled together)
  • Compression: Cross-DC traffic is compressed
  • Retry logic: Automatic retry if cross-DC write fails
  • Hinted handoff: If remote DC temporarily down, hints stored and replayed later

Monitoring Cross-DC Replication:

-- Check cross-DC replication lag
nodetool netstats

-- Example output:
Mode: NORMAL
Datacenter: US
  Outbound: 0 bytes
  Inbound: 12.5 MB
Datacenter: EU  
  Outbound: 8.2 MB (lag: 120ms)
  Inbound: 0 bytes
Datacenter: Asia
  Outbound: 6.1 MB (lag: 280ms)
  Inbound: 0 bytes

Typical Replication Lags:

  • Same region: US-East โ†’ US-West = 60-80ms
  • Cross continent: US โ†’ Europe = 120-150ms
  • Across Pacific: US โ†’ Asia = 250-300ms
  • Around globe: US โ†’ Australia = 180-220ms

Key insight: LOCAL_QUORUM achieves strong consistency locally (8ms) while eventual consistency globally (150-300ms). The client never waits for cross-DC replication, which is why it's 15-35x faster than regular QUORUM. This async replication model is what enables global-scale applications to maintain both low latency AND geographic redundancy!

5
Compare LOCAL_QUORUM vs QUORUM vs EACH_QUORUM for multi-datacenter deployments. When would you use each?

Complete Answer:

These three consistency levels represent different trade-offs between latency, consistency strength, and cross-DC behavior in multi-datacenter deployments.

Detailed Comparison (3 DCs: US, EU, Asia, RF=3 per DC):

LOCAL_QUORUM (Recommended Default):

  • Behavior: Requires QUORUM (2 of 3) from local DC only
  • Write latency: 8-12ms (local QUORUM time)
  • Read latency: 8-12ms (local QUORUM time)
  • Throughput: 30,000+ writes/sec per node
  • Consistency: Strong locally, eventual across DCs
  • DC failure impact: Writes to failed DC fail, other DCs work
  • Use case: 99% of multi-DC production workloads!
  • Examples: User profiles, orders, account settings, most business data

QUORUM (Dangerous in Multi-DC):

  • Behavior: Requires QUORUM from ANY replicas globally (with RF=3 per DC, total RF=9, needs 5 ACKs)
  • Write latency: 150-300ms (must wait for cross-DC ACKs)
  • Read latency: 150-300ms (might query remote DCs)
  • Throughput: 2,000 writes/sec (15x less than LOCAL_QUORUM!)
  • Consistency: Strong globally
  • DC failure impact: Can survive 1 DC failure (get QUORUM from remaining DCs)
  • Problem: Unacceptable latency for user-facing operations
  • Use case: Never use in multi-DC! (Use EACH_QUORUM if you need global consistency)

EACH_QUORUM (Critical Cross-DC):

  • Behavior: Requires QUORUM in EACH datacenter (2 of 3 in US AND 2 of 3 in EU AND 2 of 3 in Asia)
  • Write latency: 280-300ms (must wait for slowest DC)
  • Read latency: 8-12ms (can read from local DC only)
  • Throughput: 1,000 writes/sec (30x less than LOCAL_QUORUM!)
  • Consistency: Strong in ALL datacenters
  • DC failure impact: If ANY DC fails, all writes fail (poor availability!)
  • Use case: Critical business data that MUST be consistent globally
  • Examples: Billing records, payment processing, audit logs, compliance data

Performance Comparison Table:

Metric LOCAL_QUORUM QUORUM EACH_QUORUM
Write Latency 8-12ms โœ“ 150-300ms โŒ 280-300ms โŒ
Throughput 30k writes/sec โœ“ 2k writes/sec โŒ 1k writes/sec โŒ
DC Failure Tolerance Partial โš ๏ธ Yes โœ“ No โŒ
Global Consistency Eventual โš ๏ธ Strong โœ“ Strongest โœ“

When to Choose Each:

Choose LOCAL_QUORUM (99% of cases):

  • User-facing data: profiles, accounts, preferences
  • Business operations: orders, inventory, reservations
  • Application state: sessions, carts, user actions
  • Why: Best balance of strong consistency + low latency
  • Trade-off: Cross-DC eventual (acceptable for most use cases)

Choose EACH_QUORUM (Rare, <1% of cases):

  • Financial data: billing records, payment transactions
  • Compliance: audit logs, regulatory data
  • Critical business: must-never-lose data
  • Why: Need guarantee that ALL DCs have data before acknowledging
  • Trade-off: Much slower, but worth it for critical data

NEVER Choose QUORUM in Multi-DC:

  • It's a trap! Gives you cross-DC latency without cross-DC guarantees
  • Worse than LOCAL_QUORUM (slower) without benefits of EACH_QUORUM (global consistency)
  • If you need local consistency: use LOCAL_QUORUM
  • If you need global consistency: use EACH_QUORUM
  • There's no scenario where QUORUM makes sense in multi-DC!

Production Pattern (Hybrid Approach):

-- 99% of tables: LOCAL_QUORUM
CREATE TABLE users (...) WITH ...;
-- Use LOCAL_QUORUM for reads/writes

-- <1% of tables: EACH_QUORUM
CREATE TABLE billing_transactions (...) WITH ...;
-- Use EACH_QUORUM for writes (ensure all DCs have it)
-- Use LOCAL_QUORUM for reads (fast local reads)

Real-world example: Netflix uses LOCAL_QUORUM for 99% of their data (user accounts, viewing history, preferences) because the performance benefit is massive and cross-DC eventual consistency is fine. They use EACH_QUORUM only for critical billing data where they need absolute guarantee that payment records exist in all datacenters before acknowledging the transaction.

Key takeaway: In multi-datacenter deployments, LOCAL_QUORUM should be your default choice. Only deviate to EACH_QUORUM for truly critical data where you need guaranteed global consistency despite the huge performance cost. Never use plain QUORUM in multi-DC - it's almost always the wrong choice!

Advertisement

Responsive Ad