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.
Choose CP for payments, AP for activity feeds
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
- Shard
ledger-42has replicas inus-east-1a,1b, and1c. - A network misconfiguration isolates
1cfrom1aand1b. - The leader in
1astill has quorum (2 of 3) and continues commits. - The replica in
1ccannot reach quorum; it stops accepting writes and returns errors to any client routed there. - Load balancers health-check
1cout of the write pool within 10 seconds. - After the network heals,
1ccatches up via Raft log replication — no manual merge.
// 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.
| Component | CAP stance | Partition behavior |
|---|---|---|
| Ledger shard (per account) | CP | Reject writes without quorum |
| Balance read (post-transfer) | CP | Route to leader or quorum read |
| Transaction history API | AP-ish | Serve from lagging replica with asOf timestamp |
| Push notification queue | AP | Accept locally, dedupe on merge |
Ledger shard (per account)
CAP stanceCPPartition behaviorReject writes without quorumBalance read (post-transfer)
CAP stanceCPPartition behaviorRoute to leader or quorum readTransaction history API
CAP stanceAP-ishPartition behaviorServe from lagging replica with asOf timestampPush notification queue
CAP stanceAPPartition behaviorAccept locally, dedupe on merge
ShardPay CAP stance by component
Per-component decisions
- Ledger writes — CP by design. ShardPay will return
503during 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
replicationLagMsso 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
transferIdas 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.
System Design trackSee CAP in the system design interview framing
Quick recall
Everything you need if you only revisit this box.
- Partition tolerance is mandatory in real networks — the trade-off is C vs A during a split.
- During partition, choose C or A per component; ShardPay ledger is CP, feeds are AP.
- 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.