Consistency levels
The consistency level determines the number of replicas that must acknowledge a read or write operation as successful in order for a successful response to be returned to the client.
The consistency level defaults to ONE for all write and read operations.
You can set consistency levels for a client session, individual read/write requests, or globally.
For application development, set the consistency level with your Apache Cassandra driver.
For example, using the Java driver, call QueryBuilder.insertInto with setConsistencyLevel to set a per-insert consistency level.
With the CQL shell (cqlsh), you can use the CONSISTENCY command to set the consistency level for all queries in the current cqlsh session.
Write consistency levels
The write consistency level specifies how many replicas must respond to a write request before the write is considered successful. Even at low consistency levels, the database writes to all replicas of the partition key, including replicas in other datacenters. The write consistency level only specifies when the coordinator node can report to the client application that the write operation is considered complete. Write operations use hinted handoff to ensure the writes are completed when replicas are down or otherwise not responsive to the write request.
| Level | Description | Usage |
|---|---|---|
|
A write must be written to the commit log and memtable on all replica nodes in the cluster for that partition. |
This write consistency level provides the highest consistency, the highest latency, and the lowest availability of any level. |
|
A write must be written to the commit log and memtable on a quorum of replica nodes in each datacenter. |
Use in in multi-datacenter clusters to strictly maintain consistency at the same level in each datacenter. Writes fail if any datacenter fails to achieve a quorum. |
|
A write must be written to the commit log and memtable on a quorum of replica nodes across all datacenters. |
Use in single or multiple datacenter clusters to maintain strong consistency across the cluster. Use if you can tolerate some level of failure. For linearizable consistency, use |
|
A write must be written to the commit log and memtable on a quorum of replica nodes in the same datacenter as the coordinator node. Avoids latency of inter-datacenter communication. |
Use to maintain consistency within a single datacenter in multi-datacenter cluster. For linearizable consistency, use |
|
A write must be written to the commit log and memtable of at least one replica node. |
Provides high availability and low consistency. Sufficient for most use cases because consistency requirements are not stringent. |
|
A write must be written to the commit log and memtable of at least two replica nodes. |
Similar to |
|
A write must be written to the commit log and memtable of at least three replica nodes. |
Similar to |
|
A write must be sent to and successfully acknowledged by at least one replica node in the local datacenter. |
Achieves a consistency level of For security and compliance reasons, restricted or offline (airgapped) datacenters should always use |
|
A write must be written to at least one node. If all replica nodes for the given partition key are down, the write can still succeed after a hinted handoff has been written.
If all replica nodes are down at write time, an |
Provides low latency and a guarantee that a write never fails. Delivers the lowest consistency and highest availability. |
|
Achieves linearizable consistency by preventing unconditional updates. |
Equivalent to |
|
Same as |
Equivalent to |
Read consistency levels
For read operations, the read consistency level specifies how many replicas must respond to a read request before returning data to the client application. If a read operation detects inconsistency among replicas, the database initiates a read repair to update the inconsistent data.
| Level | Description | Usage |
|---|---|---|
|
Queries return the most recent data from all replica nodes in the cluster. All replica nodes must respond. The read operation fails if any replica does not respond. |
This read consistency level provides the highest consistency, the highest latency, and the lowest availability of any level. |
|
Queries return the most recent data from a quorum of replica nodes in each datacenter. |
Use in multiple datacenter clusters to ensure data consistency at the same level in each datacenter. Queries fail if any datacenter fails to achieve a quorum. |
|
Queries return the most recent data from a quorum of replica nodes across all datacenters. |
Ensures strong consistency across the cluster if the replication factor is high enough to tolerate down nodes or timed-out cross-datacenter communication. Cross-datacenter communication can incur extra latency. |
|
Queries returns the most recent data from a quorum of replicas in the current datacenter. Avoids latency of cross-datacenter communication. |
Use to maintain consistency within the single datacenter in multiple-datacenter clusters with a rack-aware replica placement strategy, such as NetworkTopologyStrategy, and a properly configured snitch. |
|
Queries return data from the closest replica node in the local datacenter. |
Achieves a consistency level of For security and compliance reasons, restricted or offline (airgapped) datacenters should always use |
|
Similar to Reads at |
Use to read the latest value of a column after a client has invoked an LWT to write to the column.
Reads at |
|
Same as |
Equivalent to |
|
Queries return data from the closest replica, as determined by the snitch. |
Provides the highest availability of all the levels if you can tolerate a comparatively high probability of stale data being read. The replica contacted for the read might not have the most recent writes. |
|
Queries return the most recent data from two of the closest replicas. Two replica nodes must respond. |
Assuming the replication factor is greater than 3, provides high availability if you can tolerate a comparatively high probability of stale data being read. The replicas contacted for reads might not have the most recent writes. If the replication factor is 3, this is equivalent to |
|
Queries return the most recent data from three of the closest replicas. Three replica nodes must respond. |
Assuming the replication factor is greater than 3, this level provides high availability if you can tolerate a comparatively high probability of stale data being read. The replicas contacted for reads might not have the most recent writes. If the replication factor is 3, this is equivalent to |
How QUORUM is calculated
The QUORUM level writes to the number of nodes that make up a quorum.
A quorum is calculated as 1 plus half of the sum of the replication_factor values for all datacenters (DCs), rounded down to the nearest whole number:
TRF = DC1 RF + DC2 RF + . . . + DCn RF
quorum = ( TRF / 2 ) + 1
For example, if there is one datacenter with a replication factor of 3, then a quorum is 2 nodes ((3 / 2) + 1 = 2).
The cluster can tolerate one replica down.
Examples:
-
In a single-datacenter cluster with a replication factor of 6, a quorum is 4 nodes (
( 6 / 2) + 1 = 4). The cluster can tolerate 2 replicas down. -
In a two-datacenter cluster where each datacenter has a replication factor of 3, a quorum is 4 nodes (
[(3 + 3) / 2] + 1 = 4). The cluster can tolerate 2 replica nodes down. -
In a five-datacenter cluster where two datacenters have a replication factor of 3 and three datacenters have a replication factor of 2, a quorum is 7 nodes (
[ (3 + 3 + 2 + 2 + 2) / 2] + 1 = 7).
The more datacenters, the higher number of replica nodes need to respond for a successful operation.
Similar to QUORUM, the LOCAL_QUORUM level is calculated based on the replication factor of the same datacenter as the coordinator node.
Even if the cluster has more than one datacenter, the quorum is calculated with only local replica nodes.
In EACH_QUORUM, every datacenter in the cluster must reach a quorum based on that datacenter’s replication factor for the write request to succeed.
For every datacenter in the cluster, a quorum of replica nodes must respond to the coordinator node for the write request to succeed.
Consistency level performance
Before changing the consistency level in production, use the CQL shell TRACING command to measure request performance at different consistency levels.