Clarify requirements & constraints
- Strong consistency and atomic multi-account transactions across regions
- 10k TPS global, <150ms median write for regional users
- Operational simplicity and predictable failover for finance compliance
High-level approach
- Shard ledger by account (or account-range) to localize most transactions; cross-shard (multi-account) transactions use a coordinated 2-phase commit (2PC) with a distributed transaction coordinator.
- Use Raft-family consensus for per-shard replication (readable leader, simpler operational model than Paxos). For geo-distribution, run a single Raft group per shard with replicas in multiple regions.
Consensus & leader placement
- Prefer leader-in-region close to majority of that shard’s traffic: place leader in region where shard’s primary user base is (reduces latency).
- Ensure at least 3-5 replicas across regions; configure majority quorums to include local replica plus remote ones when necessary.
Synchronous replication vs quorum writes
- Synchronous replication to majority (quorum writes) for commits — guarantees strong consistency and survives minority-region outages.
- For improved local latency, use local-leader + fast local ack + background replicate for non-critical reads (but only if business allows eventual reads). For strict writes keep synchronous quorum.
Partitioning & multi-region transactions
- Local transactions (same shard) handled by local leader and commit via Raft quorum — meets <150ms if leader colocated and network optimized.
- Cross-shard transactions: use coordinator implementing atomic commit + two-phase locking; minimize scope by encouraging single-shard operations via API design (e.g., composite accounts, routing).
Trade-offs
- Latency vs availability: majority quorum gives consistency but increases cross-region latency on leader elections/failures. Optimize with leaders colocated to demand and fast WAN (private links).
- Throughput vs complexity: more shards scale TPS but increase operational complexity and cross-shard transaction frequency.
- Operational complexity: Raft is simpler to operate; multi-region replicas require careful monitoring, automated failover, and runbooks.
Product decisions & developer experience
- Expose clear API patterns: single-shard operations for low latency, explicit multi-account transactional API with documented cost (higher latency).
- Provide SLA tiers: regional-fast (leader local, strict quorum) vs global-consistent (higher latency but cross-region atomicity).
- Measure: p95/p99 write latency, commit success rate, cross-shard transaction rate.
This balances strong consistency, predictable latency for regional users, and operational tractability for payments-grade requirements.