PrepZone Logo
PrepZone

CAP Theorem — Theory and Practice

During a partition you choose consistency or availability — ShardPay's payment path picks CP.

Why this matters

  • CAP is the most cited and most misunderstood theorem in distributed systems interviews — seniors are expected to apply it to real components, not recite definitions.
  • Partitions happen in production (misconfigured ACLs, AZ failures, cable cuts) — systems that pretend partitions cannot occur ship subtle double-spend bugs.
  • ShardPay documents a CP vs AP decision per service: the payment ledger rejects writes during quorum loss; the merchant notification feed keeps accepting events and reconciles later.
  • Interviewers probe whether you can name what you sacrifice and how you detect the partition scenario in metrics and runbooks.
During partition
ConsistencyAll nodes agree
AvailabilityEvery request gets a response
Partition toleranceNetwork splits happen

Choose CP for payments, AP for activity feeds

You cannot have all three during a network split. Pick the pair that matches your SLA.

CAP in one partition

When two ledger replicas lose connectivity, they cannot both accept writes and still guarantee that every reader sees the same balance. If both accept writes (AP), balances diverge until reconciliation. If one side rejects writes (CP), some clients see errors but committed data stays consistent across surviving nodes.

ShardPay's ledger shard runs three replicas in one region. During an inter-AZ partition, the minority side stops accepting transfers and returns 503 ShardUnavailable while the majority continues committing. Merchants see a brief error banner; no account is debited twice across divergent copies.

Key points

  • Consistency (CAP) — Every read returns the most recent write or an error. In CAP terms this is linearizable agreement on shared state, not eventual convergence. ShardPay applies this definition to the ledger primary: a balance read after a successful transfer must reflect that transfer or fail loudly.
  • Availability (CAP) — Every request to a non-failing node receives a non-error response within a bounded time. ShardPay's activity feed API stays available during partition by accepting writes locally and enqueueing them for merge — users still see "payment sent" even if cross-region sync is delayed.
  • Partition tolerance — The system continues operating despite arbitrary message loss or delay between nodes. Real WAN and cloud networks always partition eventually; pretending otherwise is not an option. ShardPay assumes partition tolerance everywhere and explicitly chooses C or A per component when a split is detected.
  • CP (Consistency + Partition tolerance) — Sacrifices availability on the minority side during a partition so writes do not fork. ShardPay's ledger is CP: if a shard cannot reach a write quorum, it refuses new debits rather than risk inconsistent balances.
  • AP (Availability + Partition tolerance) — Sacrifices strong consistency during a partition; replicas may diverge and reconcile later. ShardPay's merchant notification feed and analytics counters are AP — stale counts are acceptable; lost notifications are not.

Why partition tolerance is not optional

CAP is often misread as "pick two of three." In practice, networks partition and you still need the system to run. The real choice during a partition is between consistency and availability — not whether to tolerate partitions at all.

ShardPay's SRE runbook labels any shard where replica_reachable_count < quorum as partition mode. Dashboards switch from latency SLOs to "reject rate on minority" SLOs. On-call engineers do not try to "fix CAP" — they confirm the CP component is rejecting writes on the isolated minority and that AP components have reconciliation jobs running.

Walkthrough: AZ failure on a ledger shard

  1. Shard ledger-42 has replicas in us-east-1a, 1b, and 1c.
  2. A network misconfiguration isolates 1c from 1a and 1b.
  3. The leader in 1a still has quorum (2 of 3) and continues commits.
  4. The replica in 1c cannot reach quorum; it stops accepting writes and returns errors to any client routed there.
  5. Load balancers health-check 1c out of the write pool within 10 seconds.
  6. After the network heals, 1c catches up via Raft log replication — no manual merge.
Java
// ShardPay ledger write gate — CP behavior during quorum loss
public WriteResult commitTransfer(Transfer transfer) {
    if (!replicationGroup.hasWriteQuorum()) {
        return WriteResult.rejected("SHARD_UNAVAILABLE",
            "Write quorum lost; refusing transfer to preserve consistency");
    }
    return leader.appendAndReplicate(transfer);
}

Designing CP vs AP per component

A single product is never purely CP or AP. ShardPay splits the architecture: financial state is CP; derived views and notifications are AP. The mistake is applying one global label instead of documenting each data path.

ComponentCAP stancePartition behavior
Ledger shard (per account)CPReject writes without quorum
Balance read (post-transfer)CPRoute to leader or quorum read
Transaction history APIAP-ishServe from lagging replica with asOf timestamp
Push notification queueAPAccept locally, dedupe on merge
  • Ledger shard (per account)

    CAP stanceCP
    Partition behaviorReject writes without quorum
  • Balance read (post-transfer)

    CAP stanceCP
    Partition behaviorRoute to leader or quorum read
  • Transaction history API

    CAP stanceAP-ish
    Partition behaviorServe from lagging replica with asOf timestamp
  • Push notification queue

    CAP stanceAP
    Partition behaviorAccept locally, dedupe on merge

ShardPay CAP stance by component

Per-component decisions

  • Ledger writes — CP by design. ShardPay will return 503 during quorum loss rather than accept a debit that another partition cannot see. This protects the invariant that total debits equal total credits per shard.
  • Read models and feeds — AP with bounded staleness. Merchant dashboards may show a transfer as "pending" for 30 seconds on a lagging replica; the API exposes replicationLagMs so clients can choose to refresh from the primary.
  • Control-plane metadata — CP via Raft (see module 6). Shard routing tables must not fork — a wrong route is worse than a temporary write outage.
  • Idempotency and reconciliation — AP components still need deduplication. ShardPay's feed merger uses transferId as a natural idempotency key so duplicate events from both sides of a partition collapse to one notification.

Example: Notification feed during partition

A merchant in eu-west-1 triggers a payout while the regional link to us-east-1 is down. The EU notification service writes the event to a local outbox and returns 200 OK immediately (AP). When the link recovers, a merger job replays EU events into the global stream, skipping duplicates by transferId. The merchant may have seen the notification before the US analytics pipeline counted the payout — both are correct under their consistency contracts.

Quick recall

Everything you need if you only revisit this box.

  1. Partition tolerance is mandatory in real networks — the trade-off is C vs A during a split.
  2. During partition, choose C or A per component; ShardPay ledger is CP, feeds are AP.
  3. Detect partition mode via quorum/reachability metrics and refuse writes on the minority for CP paths.

Test yourself

Answer these before moving on — recall is what makes it stick.