PrepZone Logo
PrepZone

Replication Topologies

Leader-follower, multi-leader, and leaderless replication — when each fits ShardPay's ledger.

Why this matters

  • Topology choice drives failover behavior, write conflict resolution, and operational runbooks.
  • ShardPay uses leader-follower per shard for CP writes — the pattern most financial systems standardize on.
  • Multi-leader and leaderless topologies solve real problems (geo latency, partition tolerance) but introduce reconciliation costs ShardPay avoids on the ledger.
  • Database track articles implement these patterns in storage engines — this article frames when to pick each.
Client writes
LeaderSingle writer
Follower 1
Follower 2
Leader accepts writes; followers replicate asynchronously.

Three topologies

ShardPay's ledger shard has one leader and two followers in the primary region, plus one async read replica in a secondary region. Cross-region replica lags 100–500ms but offloads merchant history queries. Multi-leader replication was evaluated and rejected: concurrent balance updates in two regions require conflict resolution that violates ShardPay's "no negative balance" invariant without expensive locking.

Topology comparison

  • Leader-follower (primary-secondary) — Single writer (leader) accepts all mutations; followers replicate the log and serve read traffic. ShardPay's default: one Raft leader per shard handles all debits and credits; followers catch up via append-only log replication.
  • Multi-leader (multi-primary) — Multiple nodes accept writes, often per region. ShardPay's EU and US notification services use multi-leader local outboxes — conflicts merge by transferId, not by balance arithmetic.
  • Leaderless (Dynamo-style) — Any replica may accept reads/writes subject to quorum rules; no fixed leader role. ShardPay's session cache layer uses leaderless quorums for idempotency key storage where last-write-wins on TTL keys is safe.
  • Synchronous vs asynchronous replication — Sync replication waits for follower ack before commit (higher durability, higher latency). ShardPay sync-replicates within the primary region (W=3 of 3 local replicas) and async-replicates to the DR region for read scaling only — DR replica is not in the write quorum.

Leader-follower in production

Leader-follower is the workhorse for CP financial data. Failover promotes a follower; clients redirect writes to the new leader. The critical discipline is ensuring only one leader accepts writes at a time.

Walkthrough: Leader failure and promotion

  1. Shard ledger-17 leader in us-east-1a stops heartbeating.
  2. Followers in 1b and 1c detect missed heartbeat after 500ms election timeout.
  3. Follower in 1b wins Raft election (higher term, log at least as up-to-date).
  4. 1b becomes leader; 1c follows; old leader's writes rejected by term check if it recovers (fencing).
  5. ShardPay's routing service updates leader endpoint within 2 seconds via Raft metadata watch.
  6. In-flight client retries hit the new leader; idempotency keys prevent duplicate transfers.
Java
// ShardPay client — follow leader hint on redirect
public TransferResult transfer(TransferRequest req) {
    try {
        return currentLeader.execute(req);
    } catch (NotLeaderException e) {
        currentLeader = metadataService.resolveLeader(e.shardId);
        return currentLeader.execute(req);  // retry once on new leader
    }
}

When ShardPay uses other topologies

Selective topology use

  • Multi-leader for regional ingress — Payment webhooks arrive in the nearest region. ShardPay writes to a local multi-leader event log, then forwards to the CP ledger shard in the home region. Conflicts are impossible at the event level because transferId is globally unique.
  • Leaderless for ephemeral state — Idempotency keys and rate-limit counters use Cassandra-style N=3, W=2, R=2. Loss of a key after 24h TTL is acceptable; loss of a committed transfer is not.
  • Chain replication — Leader → follower → follower chain reduces write fan-out on the leader. ShardPay tested chain replication for history storage; rejected due to tail latency on the last follower during read repair.
  • Read replica isolation — Async replica in eu-west-1 never joins write quorum. ShardPay tags it topology=read-only-async in service discovery so no code path accidentally writes to it.

Example: Why multi-leader was rejected for balances

ShardPay prototyped multi-leader Postgres with logical replication in EU and US. A merchant with accounts in both regions triggered concurrent $500 debits during a partition. Each region saw sufficient balance locally; both committed; post-merge balance was -$200. Without application-level conflict resolution, multi-leader is unsafe for balance registers. ShardPay kept multi-leader only for append-only event streams.

Quick recall

Everything you need if you only revisit this box.

  1. Leader-follower is the default for CP financial data — single writer per shard.
  2. Multi-leader needs explicit conflict rules; ShardPay limits it to idempotent event ingress.
  3. Async replication creates read lag — keep async replicas out of write quorums.

Test yourself

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