Overview

Multi-master replication — also called active-active replication — means two or more database nodes each accept writes simultaneously. When the network is healthy, writes replicate between nodes and the state converges. When a network partition occurs, both nodes continue accepting writes independently. When the partition heals, the nodes have diverged and must reconcile. The reconciliation step, and the question of which write wins when two writes conflict, is what this article addresses.

SAMA’s Business Continuity Planning framework requires Tier 1 systems — systems that process payments or hold customer balances — to meet an RTO of less than 4 hours and an RPO of less than 1 hour. Active-active is the architecture that achieves an RPO of zero for synchronous replication: both data centers are always current, so a failover is instantaneous for the data layer. The alternative — single-master with a standby replica — is often the safer choice for financial data, precisely because it avoids the conflict problem entirely. The question is whether the operational simplicity of single-master is worth the RPO and RTO cost, given SAMA’s requirements.

Serializable isolation in CockroachDB and YugabyteDB resolves most conflicts automatically

CockroachDB and YugabyteDB implement serializable isolation via MVCC and Raft consensus. For most operations, the database resolves conflicts transparently: concurrent writes are detected, one is serialized before the other, and the application retries the aborted transaction. Only complex business logic that cannot be expressed as a database transaction — such as saga-style multi-entity operations — requires application-level conflict resolution. If you are using a NewSQL database with serializable isolation, the conflict resolution problem is significantly smaller than it is with asynchronous multi-master replication.

Types of Conflicts

Four categories of conflict arise when two masters accept writes concurrently:

  • Write-write conflict. Two masters update the same row concurrently. The row exists on both nodes after the partition heals with different values. This is the primary case most conflict resolution strategies address.
  • Read-write conflict. A read on one master returns a stale value (because a write from the other master has not yet replicated), and the application uses that stale value to compute an update. The update is logically incorrect even if no write-write conflict occurs. This is harder to detect because neither node records the causal dependency between the read and the write.
  • Delete-update conflict. One master deletes a row while the other master updates it. The delete-update conflict has no obviously correct resolution: was the row deleted first and the update is stale, or was the update applied to valid data and the delete is invalid? Domain semantics determine the answer.
  • Ordering conflict. Writes are applied in different orders on different masters due to replication lag, violating a causal dependency. A payment settlement that must follow a payment initiation arrives at Master 2 in reverse order — settlement before initiation. If the settlement is processed first, it may reference state that does not yet exist at Master 2.

Last Write Wins (LWW)

Last Write Wins resolves write-write conflicts by comparing timestamps: the write with the latest timestamp is the winner, and the other write is discarded. It is the simplest conflict resolution strategy to implement and the most dangerous for financial data.

The problem is clock skew. Servers in different data centers run on different hardware clocks. Even with NTP synchronisation, wall-clock timestamps on different nodes can differ by tens of milliseconds. If a write on DC1 at T=100ms is logically later than a write on DC2 at T=102ms (because the DC2 write was initiated first but arrived second due to network latency), LWW selects the DC2 write as the winner based on timestamp. The DC1 write — which was the logical successor — is silently discarded.

Hybrid Logical Clocks (HLC) address clock skew by combining physical time with a logical counter. The HLC on each node advances both when the physical clock advances and when a message is received that carries a higher logical timestamp. This creates a partial order that is consistent with causality even when physical clocks diverge. CockroachDB uses HLC internally for its MVCC timestamps.

StrategyImpl. complexityClock dependencyBanking applicabilitySAMA BCP suitability
LWW (wall clock)LowHigh — clock skew silently loses writesReference data only (idempotent)Not acceptable for financial state
CRDTMediumNone — merge-basedCounters, flags, permission setsAcceptable for non-balance state
Vector ClocksHighNone — logical clocksSibling detection; needs app resolutionAcceptable with app resolution layer
App-level resolutionHighCausal ordering requiredBalance merge, saga compensationAcceptable; requires saga design
DB-level serializableLow (app) / High (infra)HLC (CockroachDB) — bounded skewAll financial operationsPreferred for Tier 1 systems

CRDTs (Conflict-free Replicated Data Types)

CRDTs are data structures designed to merge without conflict. The merge operation is commutative, associative, and idempotent — the same set of updates applied in any order produces the same result. This means CRDTs eliminate the conflict resolution problem for the data types they model.

Four CRDT types are relevant in a banking context:

  • G-Counter (grow-only counter). Each node maintains its own counter. The global value is the sum of all node counters. Increments never conflict because each node only increments its own counter. Applicable for: login attempt counts, API call counters, event occurrence counts. Not applicable for anything that decrements.
  • PN-Counter (positive-negative counter). A pair of G-Counters: one for additions, one for subtractions. The global value is the sum of all P-counters minus the sum of all N-counters. Applicable for: inventory counts, not for account balances (because a PN-Counter cannot enforce a lower bound constraint — the balance cannot go negative).
  • OR-Set (observed-remove set). An add-wins set: if two nodes concurrently add and remove the same element, the add wins. Applicable for: permission lists, feature flag assignments, approved consumer lists. The add-wins semantics must match the business rule (it does for permissions; it may not for revocations).
  • LWW-Register (per-field LWW register). Each field has its own LWW timestamp. Concurrent updates to different fields on the same record do not conflict; concurrent updates to the same field resolve by timestamp. Applicable for: user profile data (notification preferences, display settings). Subject to the same clock skew risk as LWW.
LWW with wall clock timestamps is dangerous when clock skew exceeds your conflict window

A conflict window of 10 milliseconds and an NTP synchronisation accuracy of 50 milliseconds means LWW with wall clock timestamps will silently discard the logically correct write on every conflict where the two writes occurred within 50ms of each other. Use Hybrid Logical Clocks, or design the data type to avoid LWW entirely for any state where silent data loss is unacceptable. For financial systems, silent data loss is always unacceptable.

Vector Clocks

A vector clock is a data structure attached to every write that records the logical time at which each replica last acknowledged a write from every other replica. Each node maintains a counter that it increments on every local write. When a write is replicated, the receiving node updates its vector to include the sender’s counter.

The power of vector clocks is sibling detection: two writes are causally concurrent (siblings) if neither write’s vector clock is greater-than-or-equal-to the other’s on every dimension. When siblings are detected, the database keeps both values rather than silently discarding one, and presents both to the application for resolution. The application sees the conflict explicitly and can apply domain-specific logic to resolve it.

For a payment balance, the sibling resolution logic is: if both siblings represent debits from the same starting balance, and the combined debit exceeds the balance, reject one and compensate. If the combined debit does not exceed the balance, sum the debits and apply both. This logic is domain-specific and must be implemented by the application, but it produces a correct result rather than silently losing a write.

CockroachDB and YugabyteDB Conflict Handling

CockroachDB and YugabyteDB implement serializable isolation using MVCC with Raft consensus. Every write is a distributed transaction that requires a quorum of replicas to commit. In a multi-region deployment, this means a write on DC1 does not commit until a majority of the Raft group (which may include replicas in DC2) acknowledges it. This is synchronous replication by design — the write does not return to the application until it is durable across the quorum.

This design eliminates write-write conflicts for operations that can be expressed as database transactions. The serializable isolation level means that two concurrent transactions that conflict are serialized by the database: one commits, one aborts with a serialization error, and the application retries. The application does not see two conflicting values — it sees one committed value and a transaction error.

cockroachdb-geo-partition.sqlsql
-- CockroachDB geo-partition configuration for SAMA data residency
-- All customer financial data must reside within Saudi Arabia

-- Step 1: Create multi-region database with Saudi regions
ALTER DATABASE bank_core
  SET PRIMARY REGION = 'sa-riyadh';

ALTER DATABASE bank_core
  ADD REGION 'sa-jeddah';

-- Step 2: Geo-partition payments table by account_region
-- Rows for Saudi accounts are pinned to Saudi regions
ALTER TABLE payments
  SET LOCALITY REGIONAL BY ROW AS account_region;

-- Step 3: Verify leaseholder placement stays within Saudi Arabia
SHOW ranges FROM TABLE payments WITH details;

-- Step 4: Configure max_offset for clock uncertainty interval
-- CockroachDB uses HLC; max_offset bounds the uncertainty window
-- Lower = tighter; requires better NTP; default 500ms is conservative
SET CLUSTER SETTING kv.clock.max_offset = '250ms';

-- Step 5: Enable survival goal for zone failure (not region failure)
-- REGION FAILURE survival requires 3+ regions; for 2 Saudi DCs use ZONE
ALTER DATABASE bank_core
  SURVIVE ZONE FAILURE;

-- Step 6: Create a serializable read-modify-write pattern
-- Application retries on serialization error (SQLSTATE 40001)
BEGIN;
  SET TRANSACTION ISOLATION LEVEL SERIALIZABLE;
  SELECT balance FROM accounts WHERE account_id = $1 FOR UPDATE;
  -- business rule check: balance >= debit amount
  UPDATE accounts SET balance = balance - $2 WHERE account_id = $1;
COMMIT;
-- On SQLSTATE 40001: retry the entire transaction

Application-Level Conflict Resolution

For cases that CRDT and database-level serialization do not cover — complex saga-style operations that span multiple entities, or legacy databases that do not support distributed transactions — application-level conflict resolution is the fallback. It is more expensive to implement and more likely to contain logic errors, but it is the correct tool when the alternatives are not available.

conflict-resolution-trigger.sql (PostgreSQL)sql
-- Application-level conflict resolution trigger for multi-master PostgreSQL
-- Assumes logical replication with conflict detection plugin (pg_logical)

CREATE OR REPLACE FUNCTION resolve_balance_conflict()
RETURNS TRIGGER AS $$
DECLARE
  v_local_balance     numeric;
  v_incoming_balance  numeric;
  v_local_version     bigint;
  v_incoming_version  bigint;
  v_resolved_balance  numeric;
BEGIN
  -- NEW: the incoming row from the remote master
  -- OLD: the current local row
  v_local_balance    := OLD.balance;
  v_incoming_balance := NEW.balance;
  v_local_version    := OLD.version;
  v_incoming_version := NEW.version;

  -- If the incoming row is a pure increment/decrement relative to a common ancestor,
  -- sum the deltas and apply both. Causal order requires version vectors.
  IF v_incoming_version > v_local_version THEN
    -- Remote write is strictly newer; apply it
    RETURN NEW;
  ELSIF v_incoming_version = v_local_version THEN
    -- Concurrent write (conflict): compute resolved balance
    -- Strategy: sum partial updates, then apply business rule (no negative balance)
    v_resolved_balance := v_local_balance + (v_incoming_balance - v_local_balance);
    IF v_resolved_balance < 0 THEN
      -- Business rule violation: cannot go negative; reject incoming debit
      -- Log for manual investigation
      INSERT INTO conflict_log (account_id, local_balance, incoming_balance,
                                 local_version, incoming_version, resolution,
                                 logged_at)
      VALUES (OLD.account_id, v_local_balance, v_incoming_balance,
              v_local_version, v_incoming_version, 'REJECTED_NEGATIVE', now());
      RETURN OLD; -- keep local value
    END IF;
    NEW.balance := v_resolved_balance;
    NEW.version := v_local_version + 1;
    RETURN NEW;
  ELSE
    -- Local write is newer; discard incoming
    RETURN OLD;
  END IF;
END;
$$ LANGUAGE plpgsql;

CREATE TRIGGER balance_conflict_resolver
  BEFORE UPDATE ON accounts
  FOR EACH ROW
  WHEN (pg_trigger_depth() = 0)   -- prevent recursive trigger on self-update
  EXECUTE FUNCTION resolve_balance_conflict();
pn-counter.lua (Redis 7 CRDT-style)lua
-- Redis Lua script: PN-Counter for per-node increment/decrement
-- Key pattern: counter:{entity_id}:p:{node_id} (positive increments)
--              counter:{entity_id}:n:{node_id} (negative decrements)
-- Global value = sum(p:*) - sum(n:*)

local key     = KEYS[1]          -- e.g. "counter:account:SA123"
local node_id = ARGV[1]          -- e.g. "dc1"
local delta   = tonumber(ARGV[2]) -- positive for credit, negative for debit

if delta > 0 then
  -- Credit: increment this node's P counter
  redis.call('HINCRBYFLOAT', key, 'p:' .. node_id, delta)
else
  -- Debit: increment this node's N counter (store as positive value)
  redis.call('HINCRBYFLOAT', key, 'n:' .. node_id, -delta)
end

-- Compute global value: sum all p: fields minus sum all n: fields
local fields = redis.call('HGETALL', key)
local total_p = 0
local total_n = 0

for i = 1, #fields, 2 do
  local field = fields[i]
  local val   = tonumber(fields[i + 1]) or 0
  if field:sub(1, 2) == 'p:' then
    total_p = total_p + val
  else
    total_n = total_n + val
  end
end

-- Return global value as string (Lua returns number via redis.status_reply)
return tostring(total_p - total_n)

SAMA BCP Requirements

SAMA’s Business Continuity Planning framework classifies systems by criticality tier. Tier 1 systems — those that process IPS payments, hold customer deposits, or support core banking operations — must meet the strictest requirements:

  • RTO < 4 hours. The system must resume full operational capacity within 4 hours of a failure event. For a payment processing system, this means both data center failover and application stack recovery within that window.
  • RPO < 1 hour for most systems; RPO = 0 for financial data. SAMA’s guidance on financial data specifically requires synchronous replication: no committed transaction should be lost in a failover. An RPO of 1 hour for a payment system means up to 1 hour of payment records could be lost — which is not acceptable under any interpretation of prudent banking practice.
  • Data centre separation of ≥ 30 km. The two data centers must be geographically separated to withstand a localised disaster (power grid failure, flooding, fire). This separation also means the synchronous replication round-trip adds latency proportional to the distance — approximately 0.1ms per km of fibre at the speed of light, plus switching overhead.
  • Synchronous replication for financial data; asynchronous acceptable for analytics. SAMA distinguishes between data that, if lost, constitutes a financial loss (payment records, account balances) and data that supports operational decisions (reporting data, audit logs that are captured separately). Synchronous replication is required for the former; asynchronous replication with a documented RPO is acceptable for the latter.
  • Annual BCP test with documented results. The examination team does not accept a DR plan document as evidence. They require test results showing that the failover was executed, the RTO was met, and the data integrity of the recovery was verified. Untested BCP plans are treated as equivalent to no BCP plan.

Pitfalls

Active-active for account balances without synchronous replication is not SAMA-acceptable

An active-active deployment where both masters accept balance debits without a synchronous cross-datacenter commit is a configuration that allows the “pay twice” scenario: a customer with a 10,000 SAR balance debits 8,000 SAR on DC1 and 7,000 SAR on DC2 simultaneously during a network partition. Both checks pass locally (10,000 > 8,000; 10,000 > 7,000); both debits commit locally. The result is a negative balance of −5,000 SAR when the partition heals. This is not a data inconsistency that can be reconciled — it is a financial loss. SAMA will treat it as a regulatory breach, not a technical incident. The only safe design for account balance writes is synchronous replication with distributed transactions.

Split-brain without quorum leaves both masters accepting writes with no path to resolution

A two-node active-active deployment without a quorum mechanism (a witness node, an arbitrator, or a fencing agent) enters split-brain during a network partition: both masters believe they are the authoritative node and accept writes independently. When the partition heals, both masters have conflicting writes with no ordering information. Vector clocks require at least partial ordering to detect siblings; LWW requires clock agreement; application resolution requires a causal history. Without a quorum mechanism, none of these can be applied reliably. Always deploy an odd number of nodes with a quorum mechanism, or use a database (CockroachDB, YugabyteDB) that has Raft-based quorum built in.

Clock skew causing LWW to silently discard the correct write

This is not a theoretical failure mode. In a 2024 production incident at a multi-master deployment, a 40ms NTP drift between two data center nodes caused LWW to discard a debit transaction that arrived at the logically-later node 35ms after the conflicting credit. The credit committed; the debit was silently lost. The balance was incorrect for 14 hours before the daily reconciliation caught the discrepancy. The remediation was a full account balance audit for the affected customer pool. Use Hybrid Logical Clocks or eliminate LWW for any state where silent discard is unacceptable.

Production Checklist: Conflict-Free Active-Active Payment Service

  1. Classify every data entity by conflict sensitivity

    Before designing the replication topology, categorise every table: Can this entity tolerate LWW? (Reference data, idempotent configuration.) Does it need CRDT semantics? (Counters, flags.) Does it require distributed transactions? (Account balances, payment records, ledger entries.) The classification drives every subsequent design decision.

  2. Choose synchronous replication for Tier 1 financial entities

    Account balances and payment records must replicate synchronously. Use CockroachDB or YugabyteDB with Raft-based consensus, or configure PostgreSQL streaming replication with synchronous_commit = on and synchronous_standby_names pointing to the DR replica. Measure the added write latency (typically 1–3ms for a 30km fibre link) and size your application timeouts accordingly.

  3. Design the partition detection and quorum mechanism

    Implement a fencing mechanism that prevents both masters from accepting writes during a partition without a quorum. For CockroachDB/YugabyteDB, Raft provides this. For PostgreSQL multi-master, use Patroni with etcd or Consul as the distributed coordinator. Test the fencing under partition conditions in a staging environment before going to production.

  4. Implement the conflict detection and logging layer

    Even with synchronous replication and distributed transactions, application-layer conflicts can occur (concurrent saga steps, compensating transactions that race). Implement a conflict log table with every detected conflict, the resolution applied, and the entities involved. This table is a SAMA audit artefact.

  5. Test the partition and heal scenario end-to-end

    Use network emulation (tc netem on Linux or similar) to introduce a network partition between the two data center network segments in a staging environment. Run concurrent writes against both masters during the partition. Verify that the conflict resolution logic produces the correct result when the partition heals. Document the test and its results as a BCP test artefact.

  6. Wire reconciliation to the conflict store

    The daily reconciliation pipeline described in the companion article should include a check against the conflict log: every resolved conflict must have a corresponding reconciliation verification pass confirming that the resolved state matches the correct business state. Conflicts that cannot be automatically verified must generate an OPEN_BREAK in the break investigation table.

  7. Document and test the SAMA BCP scenario

    Produce the BCP test plan that SAMA will examine: which data center fails, which traffic is failed over, what the application failover procedure is, what the data integrity verification step is, and what the expected RTO and RPO are. Execute the test annually, record the actual RTO and RPO achieved, and file the results as a governance artefact. The test is not complete until data integrity has been verified by running the reconciliation pipeline over the failover window.