Data layer: sharding and automated failover
From 1.14: primary failover is human-speed (5–10 min) and the budget is now 4. From 0.3: ~700 cities are a natural partition key — trips, drivers, surge are all city-local. From 1.1: schemas are already private per service, so partitioning is per-service surgery, not a global rewrite.
At first glance, it seems reasonable to fail over a large global primary and let clients reconnect however they happen to retry. The choice optimizes the common case by weakening the boundary that matters. Each city shard has exactly one writable primary, promotion requires quorum and fencing, and interrupted operations remain safe to retry.
Two independent moves compose here. Sharding: the hot stores (trips, location shifts, quotes) partition by city id — a trip lives its whole life in one city, so shard-local transactions cover the invariants from 1.6 with no cross-shard coordination. A thin routing layer maps city → shard; each shard is small enough to fail over fast, back up fast, and reason about. The blast radius of any single database incident drops from "everything" to "these 40 cities" — sharding is as much an availability tool as a capacity one. Global, small datasets (users, cities, rate cards) stay unsharded on their own cluster; the ledger shards by driver/rider account, its natural key.
Automated failover: a consensus-based orchestrator (3+ observers across zones voting, so one observer's flaky network can't trigger a false failover) detects primary death, verifies the sync standby is caught up, promotes it, and repoints connection routing — target under 30 seconds. The neglected half is client behavior: every service must treat a dropped connection and a read-only error as normal, back off with jitter, and re-resolve to the new primary. Failover drills (2.8's canary spirit, 3.7's chaos) rehearse this monthly; an untested failover path is a fiction.
sequenceDiagram
autonumber
participant SVC as Trip service (client)
participant PROX as Connection router
participant ORC as Failover orchestrator (quorum of 3)
participant P as Shard 7 primary
participant S as Shard 7 sync standby
P--xORC: heartbeats stop
ORC->>ORC: 2 of 3 observers agree — primary dead, not a netsplit
ORC->>S: caught up? yes — promote
S-->>ORC: now primary
ORC->>PROX: repoint shard 7 writes to new primary
SVC->>PROX: write fails mid-transition
SVC->>SVC: retry with jittered backoff (idempotent — 1.8 keys)
SVC->>PROX: retry
PROX->>S: routed to new primary
S-->>SVC: committed
Note over P,S: total write unavailability for shard 7: ~20 s, zero humans
flowchart LR
SVC["Services"] --> RT["Shard router — city id to shard"]
RT --> S1[("Shard 1 — cities 1-40, primary + standby + replicas")]
RT --> S2[("Shard 2 — cities 41-80")]
RT --> SN[("Shard N")]
SVC --> GL[("Global cluster — users, cities, rate cards")]
SVC --> LG[("Ledger shards — by account id")]
The failover orchestrator is itself a distributed system; a naive single-observer version causes split brain — two primaries accepting writes, the worst data corruption there is. Quorum detection plus fencing (the old primary is forcibly disconnected before the new one accepts writes) are non-optional. Buy this machinery (managed databases, proven orchestrators) rather than building it.