Step 1: Why we copy data across regions
Imagine three warehouses (AWS regions) that each hold a copy of every customer rating. When one warehouse is cut off from the others, the system must still accept new ratings so users can keep watching.
What is replication factor?
Replication factor = how many copies of each piece of data must exist. In the Netflix setup the factor was set to 3, so every rating lives in three different regions.
R1 (us-east-1) ─── R2 (us-west-2) ─── R3 (eu-west-1)
▲ ▲ ▲
copy of rating copy of rating copy of rating
Step 2: How a write is accepted
Cassandra only says “write succeeded” after enough copies have acknowledged it. That “enough” number is controlled by the consistency level you choose.
What is a consistency level?
Consistency level = how many replicas must reply before the database tells the application “done.”
Step 3: LOCAL_QUORUM vs EACH_QUORUM
With replication factor 3, a quorum is 2 replicas.
- LOCAL_QUORUM
Only the 2 replicas inside the same region must reply. Cross-region copies can be late.
- EACH_QUORUM
A quorum must be reached in every region. So at least 2 replicas in each of the three regions must reply.
LOCAL_QUORUM write
R1 ──✓── R2 ──?── R3
(only R1+R2 needed)
EACH_QUORUM write
R1 ──✓── R2 ──✓── R3
(needs 2 in every region)
Step 4: What happens when the link between regions slows down
When the cable between us-east-1 and us-west-2 degrades:
- LOCAL_QUORUM writes still finish fast (only local replicas needed).
- EACH_QUORUM reads and writes now wait for the slow link, so tail latency jumps from tens of milliseconds to multiple seconds.
Netflix therefore changed both reads and writes to LOCAL_QUORUM:
// Astyanax config
ConsistencyLevel.LOCAL_QUORUM
Step 5: The hidden cost — stale data
Because the third copy may arrive late, a user can rate a title and still see the old recommendation for up to a minute.
To fix the copies eventually, Netflix added a lightweight repair job that runs every 30 seconds:
Repair job
every 30 s
│
▼
Compare and sync any missing rows
This is called anti-entropy: a background process that finds and fixes differences without blocking normal traffic.
What is anti-entropy?
Anti-entropy = background reconciliation that makes all copies identical over time, even if they temporarily disagreed.
Step 6: The unavoidable trade-off
By choosing LOCAL_QUORUM the team gained fast writes and kept sessions alive during network hiccups. They accepted a one-minute window of possibly outdated recommendations because losing a playback session hurts more than showing a slightly old row.
The expensive consistency guarantee is the one you only discover you needed after the partition already happened.
The core lessons
- Replication factor 3 means three copies; quorum needs two of them.
- LOCAL_QUORUM only counts replicas in the local region → fast but allows temporary staleness.
- EACH_QUORUM waits for every region → stronger consistency but high latency when links are slow.
- Anti-entropy jobs (30-second repair window in this case) clean up differences later.
- You always pay either in speed, in freshness, or in both; the choice only becomes visible during real network trouble.
