Chapter 14 · Consistency Levels, Quorums, Availability, and Client Guarantees
Quorum Math with Replication Factor and Why LOCAL_QUORUM Matters in Multi-DC Systems
Compute quorum requirements in multi-datacenter Cassandra and test why LOCAL_QUORUM behaves differently from global QUORUM and EACH_QUORUM.
Learning outcomes
AtlasMart expands from one region to two. A global
QUORUM copied from the single-DC design now makes
local checkout depend on remote replicas, while
EACH_QUORUM can make a remote-region outage fail
writes entirely. The goal is to derive behavior from RF per DC
instead of assuming “quorum is quorum.”
Compute global QUORUM, LOCAL_QUORUM, EACH_QUORUM, and LOCAL_ONE for a 3+3 replica topology.
Explain why global QUORUM can cross the WAN while LOCAL_QUORUM can remain DC-local.
Observe request availability when the entire remote DC is paused in an optional six-node local lab.
Use a deterministic latency simulation to illustrate WAN sensitivity without pretending localhost Docker is a WAN.
Choose DC-scoped consistency from application locality and failover requirements.
The mandatory single-DC labs use Apache Cassandra
5.0.9 in the pinned Docker image
cassandra:5.0.9, Java 17 inside the image,
cluster atlasmart-course, network
atlasmart-cassandra, nodes
atlasmart-cass-1..3, datacenter dc1,
racks rack1..rack3, 16 virtual nodes per node,
NetworkTopologyStrategy, replication factor (RF)
3, and explicit per-request consistency levels (CLs). New
tables use UnifiedCompactionStrategy (UCS), no default Time To
Live (TTL), and the Cassandra default
gc_grace_seconds unless a lesson states
otherwise. Authentication, client Transport Layer Security
(TLS), internode TLS, and remote Java Management Extensions
(JMX) are disabled only inside the isolated local learning
network. The optional application example uses Apache
Cassandra Java Driver 4.19.3. Windows learners
should use Docker Desktop/WSL-style Linux containers; commands
run inside containers unless labeled host-side. The full
two-DC lab adds three more nodes and is optional for machines
with roughly 12+ GiB free RAM and adequate disk/CPU; a
deterministic arithmetic/latency simulation is provided for
constrained machines.
docker network inspect atlasmart-cassandra >/dev/null 2>&1 || docker network create atlasmart-cassandradocker volume create atlasmart-cass-1-datadocker volume create atlasmart-cass-2-datadocker volume create atlasmart-cass-3-datadocker inspect atlasmart-cass-1 >/dev/null 2>&1 || docker run -d --name atlasmart-cass-1 --hostname atlasmart-cass-1 --network atlasmart-cassandra -e CASSANDRA_CLUSTER_NAME=atlasmart-course -e CASSANDRA_DC=dc1 -e CASSANDRA_RACK=rack1 -e CASSANDRA_ENDPOINT_SNITCH=GossipingPropertyFileSnitch -e CASSANDRA_NUM_TOKENS=16 -v atlasmart-cass-1-data:/var/lib/cassandra cassandra:5.0.9# Wait until node 1 reports UN before starting peers.docker exec atlasmart-cass-1 nodetool statusdocker inspect atlasmart-cass-2 >/dev/null 2>&1 || docker run -d --name atlasmart-cass-2 --hostname atlasmart-cass-2 --network atlasmart-cassandra -e CASSANDRA_CLUSTER_NAME=atlasmart-course -e CASSANDRA_DC=dc1 -e CASSANDRA_RACK=rack2 -e CASSANDRA_ENDPOINT_SNITCH=GossipingPropertyFileSnitch -e CASSANDRA_NUM_TOKENS=16 -e CASSANDRA_SEEDS=atlasmart-cass-1 -v atlasmart-cass-2-data:/var/lib/cassandra cassandra:5.0.9docker inspect atlasmart-cass-3 >/dev/null 2>&1 || docker run -d --name atlasmart-cass-3 --hostname atlasmart-cass-3 --network atlasmart-cassandra -e CASSANDRA_CLUSTER_NAME=atlasmart-course -e CASSANDRA_DC=dc1 -e CASSANDRA_RACK=rack3 -e CASSANDRA_ENDPOINT_SNITCH=GossipingPropertyFileSnitch -e CASSANDRA_NUM_TOKENS=16 -e CASSANDRA_SEEDS=atlasmart-cass-1 -v atlasmart-cass-3-data:/var/lib/cassandra cassandra:5.0.9# Continue only after all three nodes are UN.docker exec atlasmart-cass-1 nodetool statusdocker exec atlasmart-cass-1 nodetool versiondocker exec atlasmart-cass-1 java -versiondocker exec atlasmart-cass-1 cqlsh -e "SHOW VERSION"
Run commands only against the disposable Apache Cassandra course lab or another explicitly approved non-production environment. Confirm node, keyspace, table, container, volume, path, and datacenter targets before destructive, failure-injection, cleanup, repair, restore, security, or topology operations. Capture current state and expected rollback/recovery evidence first; output and timings can differ by host, operating system, Java runtime, Docker/runtime, driver, and Cassandra configuration.
Core terms used throughout this chapter
Replication factor (RF) is the number of
replicas Cassandra is configured to keep for each partition in a
datacenter. A replica stores a copy of a
partition; the coordinator is the node handling
one client request and is not a permanent leader. A
consistency level (CL) is the per-request rule
that says how many and which replicas must acknowledge a write
or satisfy a read before the coordinator can return success. A
datacenter (DC) is Cassandra's logical
locality/failure-domain grouping; a rack is a
smaller placement grouping within a DC. A
quorum is a majority, calculated as
floor(RF/2)+1 for the relevant replica scope.
CQL means Cassandra Query Language;
cqlsh is its shell. A
Lightweight Transaction (LWT) is a conditional
CQL operation using Paxos consensus; it has a serial phase and a
regular learn/data phase. Availability here
means whether enough replicas are alive/responding for the
requested CL—not whether the cluster has any node alive.
Staleness means a read can return an older
replica version because the requested read CL did not
necessarily intersect the replicas that acknowledged a prior
write.
1. Quorum math changes when RF spans datacenters
Suppose NetworkTopologyStrategy uses RF=3 in
dc1 and RF=3 in dc2. The total RF is
six. Global QUORUM therefore requires
floor(6/2)+1 = 4 responses from replicas across the
combined replica set. LOCAL_QUORUM in dc1 requires
floor(3/2)+1 = 2 responses from dc1 only.
EACH_QUORUM requires two in dc1 and two in
dc2. LOCAL_ONE needs one local response.
| Topology: dc1=3, dc2=3 | Required responses | Remote-DC outage |
|---|---|---|
LOCAL_ONE from dc1 |
1 in dc1 | can succeed if one dc1 replica is available |
LOCAL_QUORUM from dc1 |
2 in dc1 | can succeed with dc2 unavailable |
QUORUM |
4 of 6 globally | cannot succeed with only the three dc1 replicas alive |
EACH_QUORUM |
2 in dc1 + 2 in dc2 | cannot succeed if either DC lacks its quorum |
ALL |
6 of 6 | fails on any unavailable replica |
The familiar R + W > N intersection idea is
useful only after defining which N and which replica
scopes. Local levels intersect within a DC; global levels span
the total replica set. Lightweight transactions add a serial
phase, and failures/timeouts can make visibility more nuanced
than the slogan suggests.
2. Optional full two-DC Docker extension
These nodes join the existing cluster with
dc2 topology labels. This is still one physical
host, so it proves Cassandra's DC-aware routing/availability
arithmetic but does not reproduce real WAN
latency, packet loss, or independent power/network failure
domains.
docker volume create atlasmart-cass-4-datadocker volume create atlasmart-cass-5-datadocker volume create atlasmart-cass-6-datadocker run -d --name atlasmart-cass-4 --hostname atlasmart-cass-4 --network atlasmart-cassandra -e CASSANDRA_CLUSTER_NAME=atlasmart-course -e CASSANDRA_DC=dc2 -e CASSANDRA_RACK=rack1 -e CASSANDRA_ENDPOINT_SNITCH=GossipingPropertyFileSnitch -e CASSANDRA_NUM_TOKENS=16 -e CASSANDRA_SEEDS=atlasmart-cass-1 -v atlasmart-cass-4-data:/var/lib/cassandra cassandra:5.0.9# Wait until node 4 is UN before starting the remaining peers.docker exec atlasmart-cass-1 nodetool statusdocker run -d --name atlasmart-cass-5 --hostname atlasmart-cass-5 --network atlasmart-cassandra -e CASSANDRA_CLUSTER_NAME=atlasmart-course -e CASSANDRA_DC=dc2 -e CASSANDRA_RACK=rack2 -e CASSANDRA_ENDPOINT_SNITCH=GossipingPropertyFileSnitch -e CASSANDRA_NUM_TOKENS=16 -e CASSANDRA_SEEDS=atlasmart-cass-1 -v atlasmart-cass-5-data:/var/lib/cassandra cassandra:5.0.9docker run -d --name atlasmart-cass-6 --hostname atlasmart-cass-6 --network atlasmart-cassandra -e CASSANDRA_CLUSTER_NAME=atlasmart-course -e CASSANDRA_DC=dc2 -e CASSANDRA_RACK=rack3 -e CASSANDRA_ENDPOINT_SNITCH=GossipingPropertyFileSnitch -e CASSANDRA_NUM_TOKENS=16 -e CASSANDRA_SEEDS=atlasmart-cass-1 -v atlasmart-cass-6-data:/var/lib/cassandra cassandra:5.0.9docker exec atlasmart-cass-1 nodetool status
CREATE KEYSPACE IF NOT EXISTS atlasmart_multi_dcWITH replication = {'class':'NetworkTopologyStrategy','dc1':3,'dc2':3};CREATE TABLE IF NOT EXISTS atlasmart_multi_dc.checkout_guard ( checkout_id text PRIMARY KEY, state text, updated_at timestamp) WITH compaction = {'class':'UnifiedCompactionStrategy'};CONSISTENCY EACH_QUORUM;INSERT INTO atlasmart_multi_dc.checkout_guard VALUES ('checkout-42','OPEN',toTimestamp(now()));SELECT * FROM atlasmart_multi_dc.checkout_guard WHERE checkout_id='checkout-42';
3. Remove the remote DC from the response path
After dc2 is fully populated, pause all three dc2 containers.
Once failure detection marks them down, a dc1
LOCAL_QUORUM request still has its required two
local replicas. Global QUORUM needs four of six and
therefore cannot be satisfied by dc1's three replicas alone.
EACH_QUORUM also fails because dc2 cannot provide
its local majority.
docker pause atlasmart-cass-4 atlasmart-cass-5 atlasmart-cass-6docker exec atlasmart-cass-1 nodetool status
CONSISTENCY LOCAL_QUORUM;SELECT * FROM atlasmart_multi_dc.checkout_guard WHERE checkout_id='checkout-42';CONSISTENCY QUORUM;SELECT * FROM atlasmart_multi_dc.checkout_guard WHERE checkout_id='checkout-42';CONSISTENCY EACH_QUORUM;UPDATE atlasmart_multi_dc.checkout_guard SET state='SUBMITTED',updated_at=toTimestamp(now()) WHERE checkout_id='checkout-42';
docker unpause atlasmart-cass-4 atlasmart-cass-5 atlasmart-cass-6# Wait for UN before any cleanup.docker exec atlasmart-cass-1 nodetool status# Optional teardown after this lesson only:# docker rm -f atlasmart-cass-4 atlasmart-cass-5 atlasmart-cass-6# docker volume rm atlasmart-cass-4-data atlasmart-cass-5-data atlasmart-cass-6-data
4. Latency: simulate WAN sensitivity honestly
A single-host Docker cluster cannot produce a trustworthy cross-region latency benchmark. Use tracing to confirm which DCs participate, then model the response-order effect with explicit assumed round-trip times (RTTs). The following simulation is illustrative, not a Cassandra benchmark.
local_ms = [1.2, 1.5, 2.0]remote_ms = [75.0, 82.0, 90.0]all_rtts = sorted(local_ms + remote_ms)local_quorum_ms = sorted(local_ms)[1] # 2nd local response for RF=3quorum_ms = all_rtts[3] # 4th response for total RF=6print("LOCAL_QUORUM illustrative completion:", local_quorum_ms, "ms")print("QUORUM illustrative completion:", quorum_ms, "ms")# Replace the numbers with measured RTT/trace data from your topology.
EACH_QUORUM intentionally couples availability and tail latency to every replicated DC. That may match a global invariant, but it is a poor default when the application is DC-local and should survive remote-DC loss. Decide from the invariant, not the adjective “strong.”
Check your understanding
- With RF 3+3, what is global QUORUM?
- What is LOCAL_QUORUM in dc1?
- Can LOCAL_QUORUM in dc1 succeed when all of dc2 is unavailable?
- Can global QUORUM succeed with only the three dc1 replicas alive in a 3+3 topology?
- Why is localhost Docker not a WAN latency benchmark?
Review the answers
1. Four responses from the six-replica global set.
2. Two responses from dc1 replicas.
3. Yes, if at least two dc1 replicas are available.
4. No; it needs four responses.
5. All containers share one host/network path; it lacks real inter-region RTT, jitter, packet loss, bandwidth, routing, and independent failure domains.
Production judgment
Choose consistency from a business invariant, topology, RF, and failure budget—not from a cluster-wide slogan. Record which reads must observe which writes, whether the invariant is local to one DC or global, whether conditional uniqueness/compare-and-set is required, and what latency/availability degradation is acceptable when replicas or a whole DC are unavailable. Then test that contract with the real driver, routing policy, request timeout, retry/speculative-execution policy, idempotency classification, and representative network latency.
Consistency level does not replace durable commit-log/storage
design, repair, backup/restore, security isolation,
schema/partition design, compaction/tombstone management, or
application-level idempotency. ALL is not
“permanent durability”; LOCAL_ONE is not tenant
isolation; a timeout is not proof that a write failed; and a
successful weak write does not promise a subsequent weak read
will be fresh. Storage-Attached Indexing (SAI) or vector search
can add read work but do not redefine RF/CL arithmetic. Managed
Cassandra services may restrict topology visibility or CL
choices; verify provider semantics rather than assuming Apache
Cassandra behavior is exposed unchanged. Lesson 3 turns the math
into read-versus-write choices and demonstrates a real
stale-read window when a weak read does not intersect the
replicas that acknowledged a prior write.
Summary and next bridge
In multi-DC Cassandra, “quorum” is incomplete without scope. LOCAL_QUORUM protects local availability/latency; QUORUM can depend on remote replicas; EACH_QUORUM intentionally requires a majority in every DC. Next, combine read and write CLs and observe staleness under controlled replica divergence.
Authoritative references
Re-check these version-sensitive sources when regenerating this course.