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.
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
- Shard
ledger-17leader inus-east-1astops heartbeating. - Followers in
1band1cdetect missed heartbeat after 500ms election timeout. - Follower in
1bwins Raft election (higher term, log at least as up-to-date). 1bbecomes leader;1cfollows; old leader's writes rejected by term check if it recovers (fencing).- ShardPay's routing service updates leader endpoint within 2 seconds via Raft metadata watch.
- In-flight client retries hit the new leader; idempotency keys prevent duplicate transfers.
// 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
transferIdis 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-1never joins write quorum. ShardPay tags ittopology=read-only-asyncin 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.
Databases trackSee replication models in database engines
Quick recall
Everything you need if you only revisit this box.
- Leader-follower is the default for CP financial data — single writer per shard.
- Multi-leader needs explicit conflict rules; ShardPay limits it to idempotent event ingress.
- 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.