Data consistency
Compared to relational databases, DataStax Enterprise (DSE) doesn’t strictly adhere to atomic, consistent, isolated, and durable (ACID) transactions with rollback or locking mechanisms. By design, DSE cannot support joins or foreign keys, and, consequently, it cannot offer consistency in the ACID sense. Instead, DSE offers atomic, isolated, and durable transactions with eventual and tunable consistency that allows the user decide how strong or eventual they want each transaction’s consistency to be:
- Atomicity
-
DSE supports atomicity and isolation at the row-level, but trades transactional isolation and atomicity for high availability and fast writes. A write operation (insert, update, or delete) is atomic at the partition level, meaning batch writes to two or more rows in the same partition are treated as one write operation. DSE compares client-side timestamps (in UTC) to determine the most recent update to a column, using a last-write-wins policy to resolve conflicting responses from replicas. Even if replicas are inconsistent, reads at consistency levels above
ONEorANYare much more likely to return the most recent value for a given column due to last-write-wins resolution. - Isolation
-
DSE write and delete operations are performed with full row-level isolation. This means that a write to a row within a single partition on a single node is only visible to the client performing the operation. The operation is restricted to this scope until it is complete. All updates in a batch operation belonging to a given partition key have the same restriction. However, a batch operation is not isolated if it includes changes to more than one partition.
- Durability
-
Writes to DSE databases are durable. On the write path, replica nodes record writes in memory (memtables) and in the commit log on disk before they are acknowledged as a success. If a crash or server failure occurs before the memtables are flushed to disk, the commit log is replayed on restart to recover any lost writes. In addition to the local durability (data immediately written to disk), the replication of data on other nodes strengthens durability. You can manage the local durability to suit your needs for consistency using the
commitlog_syncproperty in thecassandra.yamlfile. Set the option to either periodic or batch. - Eventual consistency
-
Consistency refers to how up-to-date and synchronized all replicas of a row of data are at any given moment. To prioritize high availability, DataStax Enterprise (DSE) follows an eventual consistency model. Cross-replica comparisons on the read and write paths improve the accuracy of reads, and ongoing repair operations progressively reconcile inconsistencies on disk. Repairs work to decrease the variability in replica data, but constant data traffic through a widely distributed system can lead to stale data and inconsistency at any time.
Based on the CAP theorem, distributed databases like DSE are AP systems, characterized by being highly available and partition tolerant, by default. DSE can be configured to provide stronger consistency guarantees, allowing it to behave more like a CP system (consistent and partition tolerant) in situations where you need to prioritize consistency over availability. It isn’t possible to configure a distributed database into a completely CA system.
Tunable consistency
To ensure the database can provide the proper levels of consistency for its reads and writes, DSE extends the concept of eventual consistency by offering tunable consistency. You can set the consistency level as needed to achieve the desired balance between consistency, availability, and performance.
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. However, there is a tradeoff between read/write latency and consistency. Stricter consistency levels that require acknowledgments from many or all nodes provide greater accuracy at the cost of latency and cluster availability. In contrast, lower consistency levels can be less accurate, but they have the lowest latency and involve fewer nodes (leaving more nodes available to serve other requests).
- Strong consistency (immediate consistency)
-
Strong consistency prioritizes accuracy over ready/write latency. Strong consistency can be guaranteed when the condition
R + W > Nis true, where:-
R is the number of replicas required to satisfy the read consistency level
-
W is the number of replicas required to satisfy the write consistency level
-
N is the replication factor (the total number of replicas that store a piece of data)
For example, if the replication factor is 3, strong consistency is achieved by reading and writing at
QUORUM, which requires 2 of the 3 replicas, for a total of 4 replicas between any pair of read and write operations (2 + 2 > 3). For faster writes with strong consistency, writes can be reduced to 1 replica, but reads must increase to 3 replicas, resulting in slower reads.For even stronger consistency, reads or writes at
ALLensure that all replicas are involved in a given operation. -
- Eventual consistency
-
Eventual consistency prioritizes availability and lower latency over immediate accuracy. Eventual consistency occurs if the condition
R + W ⇐ Nis true. For example, if the replication factor is 3, eventual consistency is achieved by reading atONEand writing atQUORUM. A mix ofONEandQUORUMis a common configuration for eventual consistency that produces fast reads or writes with likely (but not absolute) accuracy. With eventual consistency, repair and replication mechanisms eventually apply writes to all replicas, but reads always have a chance of returning stale data.Eventual consistency is also achieved by reading and writing at
ONEbut this configuration has the highest chance of returning stale data.
Linearizable consistency
Linearizable consistency ensures immediate, serializable isolation for lightweight transactions (LWTs), also known as compare and set (CAS) transactions.
Operations like INSERT, UPDATE, and DELETE with an IF clause trigger special handling by the database to ensure linearizable consistency.
In any database, some operations must be performed sequentially without being interrupted by other operations. Examples include incrementing counters, writing unique identifiers, and updating sensitive records. These operations must be handled in a controlled manner to avoid data inconsistencies that are difficult to resolve, like duplicate user records or incorrect account balances.
However, distributed databases present a challenge when data must be read and written sequentially. Unlike a relational database that uses a balanced tree (B-tree), Apache Cassandra-based databases like DSE use a storage structure similar to a log-structured merge tree. The database avoids read-before-write operations whenever possible due to the performance penalties and complexity it introduces in a distributed system. The write path is designed to store writes in-memory, and then periodically flush them to disk in an append-only manner, avoiding the need for read-before-write operations. Repairs and compaction reconcile inconsistencies and manage the lifecycle of on-disk data, ensuring that the database maintains an accurate (with respect to eventual consistency) and efficient storage state while maximizing availability.
The eventual consistency model means that stricter consistency levels (like ALL) don’t guarantee linearizable consistency on their own, and DSE doesn’t use locking or transactional dependencies to prevent simultaneous and conflicting updates to the same record.
For example:
-
Two clients send write requests at the same time.
-
Coordinator nodes distribute the requests among the replicas independently.
-
On any given replica, the requests are applied in the order received by that particular replica. The first request overwrites the existing row, and then the second request overwrites the change from the first request. There is no guarantee that replicas will receive and apply the requests in a specific order.
This race condition results in mismatched records, making it difficult to determine which update is correct. Subsequently, identical read requests might return inconsistent results depending on the queried replicas, and repair operations might unintentionally keep the incorrect record.
Paxos protocol for linearizable consistency
When a transaction requires linearizable consistency, the Paxos protocol coordinates among replicas to ensure that the transaction is applied in a consistent and isolated manner. This isolation is similar to the serializable level that relational database management systems (RDBMS) offer. Replica data is compared, and then any out-of-date data is set to the most consistent value. In DSE, the process combines the Paxos protocol with normal read and write operations to accomplish CAS/LWT operations.
The Paxos protocol is implemented as a series of four phases. These phases are actions that take place between a proposer and acceptors. Any node can be a proposer, and multiple proposers can be operating at the same time. For simplicity, the following description uses only one proposer:
- Prepare/Promise
-
A proposer sends a message that includes a proposal number to all of the nodes responsible for storing a replica of the key. Each acceptor promises to accept the proposal if the proposal number is the highest they have received.
- Read/Results
-
After the proposer receives a promise from a quorum of acceptors, the value for the proposal is read from each acceptor and sent back to the proposer. The proposer determines which value to use for the request.
- Propose/Accept
-
The proposer then proposes the value to a quorum of the acceptors along with the proposal number. Each acceptor accepts the proposal with a certain number if the acceptor is not already promised to a proposal with a high number.
- Commit/Acknowledge
-
The value is committed to the write path. The write is acknowledged as successful if all conditions are met.
The four sets of phases require four round trips between a node proposing a lightweight transaction and any cluster replicas involved in the transaction. Additionally, the Paxos protocol runs a learn phase before writing to define which read operations will be guaranteed to complete immediately if lightweight writes are occurring. This phase uses a non-serial consistency level. All of this activity inevitably impacts performance because it generates network traffic between the nodes.
|
Due to the performance impact of implementing linearizable consistency, use LWTs only in situations where concurrency is absolutely required. |
Don’t mix LWTs with non-LWT operations
LWTs block other LWTs from occurring, but they don’t block normal read and write operations from occurring. Because LWTs use a timestamping mechanism different from normal operations, mixing LWTs and normal operations can result in errors. If you use LWTs to write to a row within a partition, you must also use LWTs for reads and other writes to that row. This applies to individual and batched operations.
For example, the following series of operations can fail because it mixes LWT and non-LWT operations:
DELETE ...
INSERT .... IF NOT EXISTS
SELECT ....
The following series of operations is valid because it uses LWTs consistently:
DELETE ... IF EXISTS
INSERT .... IF NOT EXISTS
SELECT .....
Recommended consistency levels for LWTs
The database provides specialized consistency levels for LWT operations. Depending on the scope of enforcement, the following consistency levels are recommended for serialized reads and writes:
| Use case | Read consistency level | Write consistency level |
|---|---|---|
Enforce serialization within the local datacenter |
|
|
Enforce serialization across multiple datacenters |
|
|