How does Cassandra handle inconsistencies between two replicas?

Viewed 43

I have a simple question on the strategy Cassandra opted for when the following scenario happen

Scenario

  1. At T1, replica 1 receives the write mutation like name = amit, language = english
  2. At T1 + 1, replica 2 receives the update like language = japanese where name = amit

Assume, that if the write record is not replicated on replica 2 when the update for the record has come, then how does Cassandra handles the scenario.

My Guess - May be replica 2 will check the lamport timestamp of update message say it 102 and ask replica 1 for any record which is less than 102 so that it ( replica 2 ) can execute them first then execute the update statement.

Any help would be appreciated.

2 Answers

Under the hood (for normal operations, not LWTs) both INSERTs and UPDATEs are UPSERTs - they aren't dependent on the previous state to perform update of the data. When you perform UPDATE, then Cassandra just put the corresponding value without checking if corresponding primary key exists, and that's all. And even if the earlier operations come later, Cassandra will check the "write time" to resolve the conflict.

For your case, it will go as following:

  • replica 1 receives the write, and retransmit it to other replicas in the cluster, including replica 2. If replica 2 isn't available in that moment, the mutation will be written as a hint that will be replayed when replica 2 is up.
  • replica 2 may receive new updates, and also retransmit it to other replicas.

The coordinator deals with the inconsistencies depending on the consistency level (CL) used. There are also other nuanced behaviour which again are tied to the consistency level for read and write requests.

CASE A - Failed writes

If your application uses a weak consistency of ONE or LOCAL_ONE for writes, the coordinator will (1) return a successful write response to the client/driver even if just one replica acknowledges the write, and (2) will store a hint for the replica(s) which did not respond.

When the replica(s) is back online, the coordinator will (3) replay the hint (resend the write/mutation) to the replica to keep it in sync with other replicas.

If your application uses a strong consistency of LOCAL_QUORUM or QUORUM for writes, the coordinator will (4) return a successful write response to the client/driver when the required number of replicas have acknowledged the write. If any replicas did not respond, the same hint storage in (2) and hint replay in (3) applies.

CASE B - Read with weak CL

If your application issues a read request with a CL of ONE or LOCAL_ONE, the coordinator will only ever query one replica and will return the result from that one replica.

Since the app requested the data from just one replica, the data does NOT get compared to any other replicas. This is the reason we recommend using a strong consistency level like LOCAL_QUORUM.

CASE C - Read with strong CL

For a read request with a CL of LOCAL_QUORUM against a keyspace with a local replication factor of 3, the coordinator will (5) request the data from 2 replicas. If the replicas don't match, the (6) data with the latest timestamp wins, and (7) a read-repair is triggered to repair the inconsistent replica.

For more info, see the following documents:

Related