Uber / Lyft — Match Drivers and Riders
Worked solution

Uber / Lyft — Match Drivers and Riders — a worked solution

Match a rider to the closest acceptable driver in under 3 s. Geohash, S2, surge.

Try it yourself first.

You will remember almost none of this if you read it cold. The workspace walks the same 10 stages and runs the architecture you draw through a simulator, so you find out where your version breaks before you see ours.

Open the Uber / Lyft — Match Drivers and Riders workspace

The problem

Build the matching/dispatch core of a ride-hailing marketplace at Uber/Lyft scale. Drivers stream their GPS over a persistent socket; riders tap "request" and expect a driver's phone to ring within ~3 seconds. The system must pick a driver who is geographically close, currently available, likely to accept, and quote a fare with a surge multiplier that's both real-time and stable enough not to flicker between offer and confirm. Out of scope for this canonical: payments, in-trip telemetry, fraud, accounting — this is the matching main flow only.

The interesting pressure is: (a) the geo-index is write-dominant at ~125:1 vs reads — every driver heartbeat is a write, only one rider request consumes ~19 cells of fan-out — flipping every assumption the typical "read-through cache" canonical relies on; (b) demand is violently skewed — a single H3 r9 cell at MSG when a concert lets out can take 100x normal load while everything else is idle; (c) the dispatch path must be sharded by geography, not by tenant, because that's the only way one CPU can hold authoritative in-memory state for a contiguous patch of map.

The reference architecture

Reference architecture for Uber / Lyft — Match Drivers and Riders: 20 components — Rider Client, Driver Mobile, Load Balancer · Global Edge, API Gateway · Public Edge, Service · Ride Request, Service · ETA & Pricing, Service · DISCO Dispatcher, Service · Location Ingest, Coordinator · RingPop, KV Store · Driver Live Index, Cache · Surge Multipliers, NoSQL DB · Rides & Offers, Stream · Marketplace Bus, Worker · Surge Aggregator, External · Push (APNs / FCM), External · Maps & Routing, Monitoring · M3 + Prom + Alertmanager, Scheduler · Offer Reaper, Worker · Outbox Relay (CDC), Queue · Poisoned Offer DLQ — connected by 31 flows.Rider ClientclientDriver MobilemobileLoad Balancer · Globa…Anycast L4 + GeoDNS pin…API Gateway · Public …Envoy 1.28 + WAF + JWT …Service · Ride RequestGo 1.22, gRPC, idempote…Service · ETA & Prici…Go + Python (DeepETA mo…Service · DISCO Dispa…Go 1.22; in-memory cell…Service · Location In…Go; pipelined Redis cli…Coordinator · RingPopRingPop (SWIM gossip + …KV Store · Driver Liv…Redis 7 cluster mode, 6…Cache · Surge Multipl…Redis 7 cluster, 16 sha…NoSQL DB · Rides & Of…Cassandra 4.1, RF=3, QU…Stream · Marketplace …Kafka 3.7, RF=3, 24 h r…Worker · Surge Aggreg…Apache Flink 1.18, Rock…External · Push (APNs…Apple APNs HTTP/2 + Goo…External · Maps & Rou…Google Maps Routes API …Monitoring · M3 + Pro…M3 (Uber) for metrics, …Scheduler · Offer Rea…Go cron, leader-elected…Worker · Outbox Relay…Debezium-style CDC tail…Queue · Poisoned Offe…Kafka topic dispatch.ev…
20 components, 31 flows. A dashed line is an asynchronous hop. This is the reference design, not the only one that works.

What each component is for

Rider Client

Issues GET /eta as the user drags the pin (browse) and POST /requests on tap-to-confirm. Holds a persistent WebSocket through the gateway for live offer status (searching → matched → driver_arriving). Re-uses a client-generated request_nonce on every retry so the matcher dedupes.

Why it exists. We considered a stateless polling client (poll /requests/:id every 500 ms) — at 5K peak rider QPS x 4 polls/s that adds 20K spurious QPS to the gateway with no UX benefit. WS is the load-bearing choice; polling is the fallback when WS keepalive fails twice.

When it fails. App backgrounded mid-search drops the WS; searching state goes silent. Detection: client telemetry on ws_reconnect_rate and request_to_match_seconds p99. Mitigation: server-side reaper expires the request after offerTtlSec x 4; reconnect carries the same request_nonce so the matcher returns the cached in-flight result instead of a duplicate dispatch.

Driver Mobile

Streams GPS samples over WebSocket every 4 s while idle, every 1–2 s with a passenger or above 50 km/h. Receives offer envelopes via push (APNs/FCM) AND the same WS; POSTs accept/decline within the offer TTL. Snaps locations to the road graph client-side before send, so the server doesn't see oscillation around a parking lot.

Why it exists. Driver app is a bidirectional system citizen — both the heaviest write source (one heartbeat per driver every 4 s) and the receiver of the dispatch fanout. We rejected a polling-only design: at 1.5 M concurrent drivers, the reconnect storm on every gateway deploy would be ~375 K new TLS handshakes, which crushes the LB. Persistent WS with graceful drain on the server side absorbs it.

When it fails. GPS drift puts the driver in the wrong H3 cell — rider gets a driver-across-the-river match. Detection: server-side driver_jump_distance_m_per_s > 60 (> 220 km/h) is treated as suspect and the location update is dropped. Mitigation: Kalman smoothing + map-matching at ingest (locingest); legitimate jumps (tunnel, ferry) tolerate one cycle of stale before reverting. Driver-tunnel WS-drop while holding an OFFERED offer: WS-keepalive miss for >2 s triggers a fast-EXPIRE on the offer (DISCO transitions OFFERED→EXPIRED, riderapi rolls forward to next candidate).

Load Balancer · Global EdgeAnycast L4 + GeoDNS pin per city_id

Terminates the rider's TLS at the closest PoP, performs a city_id lookup against a small in-memory pin table (maintained out-of-band by ops; ~10K cities globally), and proxies to the gateway in the city's home region. Holds the global view 'NYC always lands in us-east-1 dispatch, SF always in us-west-2'. WS upgrades are stickied to the same region for the connection lifetime.

Why it exists. Without an explicit region-pin layer, our 'active-active per-city' claim is aspirational — a rider whose DNS resolves to us-east-1 but whose city is pinned to us-west-2 would dispatch against an empty driver index. The alternative (let each region's gateway forward across regions) doubles cross-region traffic on every request and adds 30–80 ms median; that's the entire ETA budget. Pinning at the edge is the only place this is cheap.

When it fails. GeoDNS resolves a rider in a city whose pin was just rotated mid-deploy → rider's request lands in the old region → returns 'no drivers near you' from the stale geoidx. Detection: per-city wrong_region_request_rate > 0 (the gateway rejects with a hint header). Mitigation: gateway returns 307 Temporary Redirect with the correct region's anycast endpoint; rider auto-retries; rotation procedure includes a 60s overlap window where both regions accept reads.

API Gateway · Public EdgeEnvoy 1.28 + WAF + JWT verifier

Inside one region. Validates the rider/driver JWT, enforces per-(IP, user_id, city_id, h3_r7_cell) rate limits, upgrades to WebSocket for the heartbeat and offer streams, and routes by URI: POST /requests → riderapi, GET /eta → riderapi, POST /offers/:id/accept → riderapi, WS /v1/driver/heartbeat → locingest. POST /requests uses ring-hash LB on the Idempotency-Key header so a retried request lands on the riderapi replica already holding its nonce future; other routes round-robin. Envoy outlier-ejection sheds slow upstreams without killing the whole pool.

Why it exists. Lyft's Envoy mesh exists exactly for this — uniform retry budgets, circuit breakers, mTLS upstream. We considered terminating WS at riderapi/locingest directly; that loses one place to enforce city-level admission control, which is the single most useful lever during a hot-cell event. The gateway is where 'shed before it cascades' happens.

When it fails. LB roll without drain: 1.5M drivers reconnect in a 5 s window, TLS handshake p99 spikes past 2 s, ingest brownouts city-wide. Detection: driver_ws_reconnect_rate > 10x baseline AND tls_handshake_p99 > 2s for 60 s. Mitigation: graceful drain + server-issued Retry-After 0–30 s uniform; pre-warmed Envoy capacity at +20% headroom over peak; emergency runbook flips the 5% rolling cap down to 1% during recovery.

Service · Ride RequestGo 1.22, gRPC, idempotency-key-aware

Owns the rider-side request lifecycle. On POST /requests: validate, fan out to ETA service for the offer envelope, then call DISCO with (rider_id, request_id, lat, lng, surge_version, eta_envelope). On POST /offers/:id/accept from a driver: validate offer is still in OFFERED state, route to DISCO for the CAS-commit. Persists the request row (SEARCHING → MATCHED | NO_DRIVER | CANCELLED) to Cassandra so poll/cancel resolve from any replica. Stateless except for an in-memory LRU of request_nonce → in_flight_promise so a retried request joins the same future — the gateway ring-hashes POST /requests on Idempotency-Key, so retries land on this same replica.

Why it exists. We considered making DISCO public (no orchestrator) — DISCO's match logic must not block on ETA composition or external Maps calls; coupling them turns one Maps blip into dispatch downtime. The orchestrator is where ETA + match are fanned out concurrently, the offer envelope is built, and idempotency is enforced before DISCO ever sees a duplicate.

When it fails. ETA service degrades; riderapi blocks waiting on the offer envelope; rider retries; pile-on at riderapi. Detection: riderapi_eta_call_p99 > 400 ms AND inflight_concurrent > 80% of limit. Mitigation: ETA timeout is 500 ms hard with circuit-breaker (no retry doubles the budget); on open, riderapi proceeds with surge=1.0 — degraded UX (price might shift on confirm) is better than 'searching forever'.

Service · ETA & PricingGo + Python (DeepETA model server)

Two endpoints: GET /eta (browse) and GET /offer-envelope (called by riderapi). Reads the per-cell surge multiplier from the surge cache, asks the Maps service for route distance/duration, applies the DeepETA residual model (a few-ms encoder-decoder on top of the routing ETA), and returns {eta_seconds, base_price, surge_multiplier, surge_version, ttl_seconds}. Stateless — model weights pulled on startup, refreshed via canary rollout.

Why it exists. Maps and surge can't be on DISCO's critical path: they are the long-tail latency sources. ETA is a separate process with its own SLO (200 ms p99) so a slow Maps call cannot push the matcher past its 600 ms budget. We considered baking ETA into riderapi; that fuses the fast-fail circuit breaker with the matcher's, which means a Maps outage takes down dispatch — unacceptable.

When it fails. DeepETA model rolls bad; ETA p99 jumps to 600 ms; riderapi hits its 500 ms timeout on every call. Detection: eta_response_p99 > 250 ms for 90 s + canary error budget. Mitigation: progressive rollout 1% → 10% → 100% per region with auto-rollback on p99 regression > 50 ms; fallback path returns routing ETA without the residual.

Service · DISCO DispatcherGo 1.22; in-memory cell index per shard; offer state machine

Receives match requests for a given H3 r9 cell, expands kRing(k=1..3, adaptive), pulls candidate driver_ids from the in-memory cell index plus a fast geoidx fallback, scores by ETA × acceptance-rate × fairness penalty, picks the best, CAS-claims the driver in geoidx, INSERTs the offer into rides Cassandra, fires push to APNs/FCM, and emits a transactional outbox event. State machine: OFFERED → ACCEPTED | DECLINED | EXPIRED, plus a fast-EXPIRE transition when the driver's WS goes silent for >2 s while holding an offer (avoids the rider waiting the full TTL).

Why it exists. This is the core of the product — the answer to 'who picks up this rider'. Sharded by RingPop on H3 r8 supercell (covers ~10 r9 cells, ~1 km^2) so a contiguous patch of map collapses onto one process. Shard-local in-memory state means kRing iteration is microseconds, not network calls. Alternative: every match call hits a central matcher reading from Redis on every kRing — the 4 RTTs per request blow the 600 ms budget at 5K req/s.

When it fails. RingPop ring instability mid-deploy: two owners both think they own cell C, both dispatch to driver D, driver gets two offers. Detection: offers_routed_to_wrong_owner > 0 (the rides CAS counter). Mitigation: lease-based ownership with monotonic epoch; rides reject stale-epoch writes — the second offer fails closed, riderapi falls back to the next candidate; ring re-converges within 10 s gossip cycles. Driver-WS-disconnect mid-OFFER: WS keepalive miss → fast-EXPIRE → next candidate within 2 s instead of 15 s wait.

Service · Location IngestGo; pipelined Redis client; Kafka producer; map-match filter

Terminates the driver WS heartbeat stream, runs a server-side map-match + Kalman filter on each sample (rejects > 60 m/s jumps as suspect), pipelines a SET driver:{id} + ZADD cell:{h3_r9} pair into geoidx (key=driver_id, value=(h3_cell, lat, lng, ts, status)), and publishes the same sample to the Kafka marketplace bus for surge + analytics consumers. No request fanout, no synchronous fan-in — fastest possible path from socket to KV.

Why it exists. Ingest is the write-dominant tier (1:125 read:write at the geo-index) and absolutely cannot share a process with the matcher — a slow matcher would back-pressure the heartbeat path and dispatch goes blind. Separate process means we can size locingest for write throughput (no big TLS handshake budget, fat batched Kafka producer) while DISCO sizes for matcher CPU.

When it fails. Map-match Kalman filter is too aggressive after a model bump: legitimate tunnel exits get dropped, cars vanish from the index for 30 s, riders in that area get 'no cars available'. Detection: dropped_samples_rate > 2% city-wide AND geoidx_active_drivers step-down. Mitigation: filter is feature-flagged per region; rollback on the regression alert; heartbeat samples are still emitted to Kafka raw so we can replay against a fixed filter.

Coordinator · RingPopRingPop (SWIM gossip + FarmHash CH ring), in-process

Maintains the consistent-hash ring of DISCO members and the H3-supercell → owner mapping. Every DISCO process embeds a RingPop client; gossip-converged membership lets any process forward a request to the current owner of a given cell. Holds the per-shard monotonic epoch that the dispatcher stamps onto every offer. The 'edge' from disco to ring is in-process (microseconds) — modeled separately on the diagram for clarity, not because there's a network hop.

Why it exists. Without an embedded ring, DISCO would need a centralized router (etcd / ZK on every match call) — that's a coordination hop on the 600 ms critical path, plus a second SPOF. RingPop's SWIM gossip is eventually consistent membership without a coordinator, paying the cost in bounded inconsistency during deploys (which we contain with the epoch). We considered Raft for membership; convergence latency at a 200-node ring is tighter than gossip, but Raft adds a leader and Raft elections during partial partitions are precisely the failure mode that paged on-call last incident — gossip fails open, Raft fails closed.

When it fails. Gossip storm during a 200-node deploy: membership views diverge for 20 s, cells oscillate ownership, in-flight offers route to the wrong owner. Detection: ringpop_membership_changes_rate > 5/min for 60 s + offers_routed_to_wrong_owner > 0. Mitigation: rolling deploy capped at 5% concurrency; epoch on every offer means the second 'owner' commits with a stale epoch and the rides CAS rejects it; double-dispatch is averted, the rider falls to the next candidate.

KV Store · Driver Live IndexRedis 7 cluster mode, 64 shards, 3-AZ, 2 read replicas/shard

Two key spaces: (1) driver:{id} → (h3_r9, lat, lng, status, ts, epoch) for the live driver record, TTL 30 s; (2) cell:{h3_r9} → SortedSet of driver_id scored by ts, refreshed on every heartbeat. DISCO does ZRANGEBYSCORE cell:{c} now-30s now × kRing(k) for candidate retrieval (allowed to read a follower replica, eventual), then a scripted CAS (WATCH driver:{id}; status==available; SET status=offered, epoch={E}) for the claim — the CAS routes to the leader for that shard, linearizable per record.

Why it exists. Hot working set is ~96 MB total (1.5M drivers × 64 B). A B-tree DB on disk would crush the cache hit rate at 1.875M writes/s peak. We considered Cassandra here — Schemaless is what Uber ran originally — but heartbeat write amp and the every-4-s churn turn the SSTable compaction cost into a CPU sink that costs more than just keeping it in RAM. Cassandra survives as the durable canonical (rides); Redis is the hot path.

When it fails. Async-replicated leader failover loses up to 5 s of recent heartbeats; followers promote with a stale view. Detection: cluster-mode +failover-end event paired with geoidx_lag_seconds > 5 (cluster-node-timeout default ~15 s). Mitigation: heartbeats refill the index in <2 cycles (8 s); during the gap, kRing reads degrade to LOCAL_ONE consistency, and the offer commit's CAS-on-driver-status fails closed when a driver record is missing — the matcher falls back to the next candidate rather than dispatching to a ghost. Hot-cell ZADD storm (concert end, ~2K ZADDs/s onto one shard): the cell sortedset is replicated across 2 followers so the kRing-read load shifts to followers, leaving the leader free for ZADDs.

Cache · Surge MultipliersRedis 7 cluster, 16 shards

Stores surge:{h3_r7} → (multiplier, version, computed_at) for every active res-7 supercell (~5 km^2, ~6M keys cardinality, hot subset ~500K). Read on every ETA call (~16K QPS peak — the 10:1 browse fan-out is already baked into that) — single GET per offer. Written by the surge aggregator at the end of each sliding window.

Why it exists. Surge has to be computable in milliseconds on every request — the surge multiplier is part of the offer envelope quoted to the rider. We rejected reading directly from Pinot (the analytics tier) on every request: Pinot p99 is 50 ms, single-digit ms is what the surge slot in the dispatch budget allows. The cache is the load-bearing decoupler between the windowed-aggregator latency (minutes) and the dispatch latency (milliseconds).

When it fails. Cache cold after a region-wide Redis restart: every ETA call falls back to surge=1.0 default for ~30 s while the aggregator backfills. Detection: surge_cache_hit_rate < 95% for 60 s. Mitigation: ETA service is allowed to serve surge=1.0 with a stale_surge=true flag rather than block on a slow refill; aggregator has a bootstrap-from-Kafka startup mode that reads the last 5 windows in parallel, restoring the hot subset of ~500K keys in under 30 s.

NoSQL DB · Rides & OffersCassandra 4.1, RF=3, QUORUM, 256 vnodes/host

Two tables: offer(offer_id PK, ride_id, driver_id, rider_id, state, surge_version, epoch, expires_at, ...) and ride(ride_id PK, driver_id, rider_id, status, pickup, drop, fare, started_at, ...). Idempotent writes keyed on offer_id (UUIDv4 from DISCO) so a retried INSERT is a no-op. Reads are point lookups on offer_id for state checks; the offer reaper does a WHERE expires_at < now() time-bucketed scan over a separate offer_by_expiry(time_bucket, offer_id) table to avoid SELECT * tombstone walks.

Why it exists. Cassandra (or Schemaless on top of MySQL — Uber's actual canonical) was chosen for marketplace state because the workload is write-heavy (5K offer writes/s peak, mostly idempotent INSERTs), single-key ACID is sufficient, and we need cross-region replication that doesn't block on a leader. Postgres-HA was rejected: even with sharded Citus, a write spike on one shard during an event sets up a queueing pile-up that the QUORUM-leaderless model absorbs through hinted handoff.

When it fails. Quorum lost in one AZ during a network partition: W=2 across the surviving 2 AZs still works but read repair lags; if a SECOND AZ blips, both W and R fail closed and DISCO sees quorum_unavailable. Detection: cassandra_quorum_unavailable_rate > 0 keyspace-scoped. Mitigation: explicit 'degraded dispatch' mode lowers reads to LOCAL_ONE (kRing-side reads only) but keeps writes at QUORUM — accept staleness on candidate search, never on offer commit; widen kRing radius and lengthen offer TTL so the matcher tries again on a fresh quorum.

Stream · Marketplace BusKafka 3.7, RF=3, 24 h retention, per-topic acks tuning

Three topics with different durability tiers: (1) loc.heartbeats — every driver GPS sample, ~1.875M msg/s peak, 256 partitions by city_id, producer acks=1 (heartbeats are best-effort; dropping a few is fine because the next sample replaces them in 4 s); (2) dispatch.events — offer state transitions (OFFERED, ACCEPTED, EXPIRED) for downstream surge + analytics + audit, ~5K msg/s, 64 partitions by offer_id, producer acks=all + min.insync.replicas=2 (offer events are the outbox — RPO 0 in-region, that's load-bearing); (3) surge.outputs — windowed multipliers for cache backfill + audit, ~50K writes/s peak, 32 partitions, acks=1. Cross-region replicated via uReplicator at <5 s lag (RPO bound for the dispatch outbox).

Why it exists. The bus is the load-bearing decoupler between the synchronous dispatch path (sub-second) and the windowed surge / analytics path (minutes). It also serves as the durability layer for driver heartbeats (Redis is non-persistent — durability lives here on a 24 h retention) and as the transactional outbox for offer commits (DISCO writes to Cassandra and the bus in the same offer-commit group; if the bus enqueue fails, the outbox-relay worker re-drives from Cassandra CDC).

When it fails. Consumer lag on loc.heartbeats: surge aggregator computes off stale supply, multiplier flickers vs offered price. Detection: consumer_lag_seconds{topic=loc.heartbeats} > 30 for 60 s. Mitigation: surge offer envelope is version-stamped; offers honor the quoted price for offerTtlSec on confirm regardless of refresh. Aggregator backlog drains by horizontal-scaling consumers; partition count chosen so a 2x consumer scale-out is always possible without re-keying. Poisoned-message handling: bus → DLQ (Queue · Poisoned Offer DLQ) routes events that fail consumer parse 3 times to the DLQ; on-call replay tool drains DLQ after fixing the consumer.

Worker · Surge AggregatorApache Flink 1.18, RocksDB state, Beam SDK

Streams loc.heartbeats and dispatch.events from the bus, keyed by h3_r7_cell. Maintains a 1-minute sliding window with 10-second slide. At each slide: counts unique drivers (supply) and rider requests (demand), divides into a ratio, smooths with EWMA over the trailing 3 windows, looks up the cell-specific surge curve, and writes (cell, multiplier, version, computed_at) to the surge cache + the surge.outputs audit topic. Poison messages are routed to the DLQ after 3 parse failures, never blocking the stream.

Why it exists. The surge multiplier is the only piece of dispatch metadata that needs windowed aggregation across all of supply and demand — every other read (driver location, ride state) is a point lookup. Flink owns this because (a) exactly-once semantics matter (double-counting demand inflates surge), (b) RocksDB-backed keyed state at 6M keys × 3 windows is tractable, (c) checkpoint-based recovery means a crash replays from the last consistent window, not from scratch. Spark Streaming was rejected: micro-batch latency floor is 1–2 windows, kills the 'react to a stadium ending' story.

When it fails. RocksDB state grows past TM heap during a deploy + traffic spike: the aggregator stalls; surge multipliers stop refreshing; ETA falls back to the last cached value plus the 60 s TTL window. Detection: flink_state_size_bytes > 80% TM heap AND checkpoint_duration_p99 > 60 s. Mitigation: state TTL on inactive cells (reaped after 5 unobserved windows), checkpoint to GCS not local disk so a TM swap doesn't lose state, manual repartition (re-key by region) when a single TM holds > 1M cells.

External · Push (APNs / FCM)Apple APNs HTTP/2 + Google FCM HTTP v1

Receives a (driver_token, offer_id, expires_at, surge_version) envelope from DISCO and delivers a high-priority push to the driver's iOS or Android device, waking the app to display the offer. APNs/FCM is the wakeup, NOT the transport — the actual offer payload is delivered via the WS the driver app already holds; the push only cues 'there's an offer waiting on the wire'.

Why it exists. iOS aggressively kills backgrounded sockets after ~30 s; without push, half of off-duty drivers would never get an offer. We considered making push the transport itself (envelope in the push payload) — that loses guaranteed-delivery: APNs may drop a push under load, and we'd need an out-of-band ack. Push-as-wakeup + WS-as-transport is the load-bearing belt-and-suspenders.

When it fails. APNs returns 503 for 10 minutes (an actual recurring incident pattern). Detection: push_provider_5xx_rate > 5% for 60 s, paired with offer_no_driver_ack_rate rising. Mitigation: WS-only fallback for foregrounded drivers (most of the active pool) — backgrounded iOS drivers go silent, rider sees longer match time and the matcher widens kRing to compensate. Rider-side UX shifts to 'matching is taking longer' rather than failing closed.

External · Maps & RoutingGoogle Maps Routes API + in-house OSRM fallback

Returns a routing result (distance_meters, duration_seconds, route_polyline) for an (origin, destination) pair. Used by ETA service to compute the base ETA before the DeepETA residual is applied, and by the Driver app for in-trip navigation (out of scope for matching). Treated strictly as a third-party with bounded availability and rate limits.

Why it exists. Building a worldwide routing engine is a separate billion-dollar problem. We pay Google's per-call price for the SLA they guarantee and run an in-house OSRM cluster as a degraded fallback for when Google's API rate-limits us or returns 5xx. Doing matching without routing is theoretically possible (haversine straight-line distance × fudge factor) but the ETA accuracy hit costs us match quality.

When it fails. Google Maps Routes API returns 429 (rate-limited) during a regional traffic spike. Detection: maps_429_rate > 1%. Mitigation: ETA service auto-flips to the OSRM fallback for that region, doubles the per-request timeout (OSRM is slower), and tags the offer envelope with routing=osrm for observability. Riders see a few-second ETA delta (worse on traffic-heavy routes) but the matching path itself never blocks on Maps.

Monitoring · M3 + Prom + AlertmanagerM3 (Uber) for metrics, Prometheus federation for tactical, PagerDuty

Scrapes /metrics from every service (gateway, riderapi, eta, disco, locingest, surgeagg, ring, glb, reaper, outboxrelay) and a sidecar exporter for every stateful tier (geoidx, surge, rides, bus, dlq). Aggregates per-(region, city, h3_r7_cell) so the on-call queries dashboards by cell, not by host. Owns the SLO definitions and the alert routes (page on dispatch availability < 99.95% over 5 min, p99 > 3 s sustained, heartbeat-to-visible > 8 s, offer_acceptance_rate drop > 10pp vs trailing 24h, surge_version_mismatch_on_confirm > 0.5%, offers_routed_to_wrong_owner > 0).

Why it exists. Without an explicit observability tier on the diagram, the SLOs and detection signals scattered across every node's failureMode have no system that actually measures and alerts on them. We considered embedding observability into each service (Datadog agent / OpenTelemetry collector per pod); that's the implementation, but the canonical needs to surface the dependency surface — what gets monitored, what pages, who owns the rotation.

When it fails. M3 ingestion buffer fills during a metrics tsunami (deploy + spike): metrics arrive late, alerts fire 2–5 min after the user impact starts. Detection: m3_ingest_lag_seconds > 60. Mitigation: pre-aggregation at the exporter level (per-host roll-ups) so a 1000-host fleet emits one rollup per service per minute, not 1000; alert rules use 5-min windows so a 2-min ingest blip doesn't suppress detection; tactical Prometheus federation continues to fire on tactical thresholds independently.

Scheduler · Offer ReaperGo cron, leader-elected via etcd lease, 1 active per region

Every 1 second scans the offer_by_expiry(time_bucket, offer_id) Cassandra table for offers whose expires_at < now() AND state = OFFERED, force-transitions them to EXPIRED via a CAS UPDATE on rides, and publishes EXPIRED events to dispatch.events on the bus so DISCO + downstream consumers see the transition. The reaper exists outside DISCO precisely so a poisoned in-memory state machine inside DISCO cannot wedge a rider waiting on a stuck offer.

Why it exists. DISCO drives the offer state machine in-memory; if a DISCO process crashes mid-OFFER (or its state machine wedges on a buggy code path), the offer is stuck OFFERED forever from the rider's perspective. We considered making the rider-side WS reconnect kick the reaper; that couples liveness to the rider being online, which is exactly the case where the rider is most frustrated. External sweeper is the load-bearing liveness guarantee.

When it fails. Reaper itself crashes: stuck offers sit OFFERED until a backup leader picks up the etcd lease (sub-10 s). Detection: scheduler_leader_election_lag_seconds > 10 OR offers_in_OFFERED_age_p99_seconds > 30. Mitigation: 3 standby reapers per region with 5 s lease renewal; alert if zero healthy reapers for >30 s; runbook includes a manual-EXPIRE tool that ops can run from a bastion if all reapers fail.

Worker · Outbox Relay (CDC)Debezium-style CDC tailing Cassandra commit log → Kafka

Tails the Cassandra commit log via a CDC connector for the offer table, identifies rows whose relayed_at is null beyond a 5-second threshold (i.e., DISCO's primary path failed to enqueue to bus), republishes the event to dispatch.events, and stamps relayed_at. The CDC stream is the recovery path that makes the offer-commit outbox correct under partial Kafka failure.

Why it exists. DISCO's primary path is: write Cassandra (durable) → publish bus (best-effort relay). If Kafka is briefly unavailable when DISCO tries to publish, the offer record is durable but the event is lost — downstream consumers (surgeagg, analytics, audit) miss the OFFER. We considered making DISCO retry Kafka indefinitely on failure; that blocks the dispatch return, busting the SLO. The CDC relay decouples durability from delivery.

When it fails. CDC connector falls behind Cassandra write rate (5K offer/s peak): unrelayed events accumulate; downstream consumers see growing OFFER lag. Detection: outbox_relay_lag_seconds > 30 OR cdc_tailer_offset_behind_writes > 10000. Mitigation: 8 connector instances partitioned by hash(offer_id) so a single-connector slowdown doesn't block the rest; alerting fires before lag reaches the 24h Kafka retention floor; operator runbook drains the lag by spinning additional connectors temporarily.

Queue · Poisoned Offer DLQKafka topic dispatch.events.dlq, RF=3, 7d retention, manual drain

Receives offer events (or surge windowed inputs) that surgeagg or other consumers fail to parse 3 consecutive times. Holds them for 7 days on a separate Kafka topic so the on-call has time to identify the schema regression, fix the consumer, and replay. Never auto-replayed into the hot path — replay is an explicit ops action.

Why it exists. Without a DLQ, a poisoned message (e.g. a surge_version that no longer parses after a model rollback) cycles forever in surgeagg's fixed-retries loop. At 5K offer/s a 1% poisoned rate is 50/s into an unbounded retry — within minutes the consumer is pinned on bad messages, lagging the legitimate ones. DLQ gives consumers an exit valve.

When it fails. DLQ depth grows during a sustained schema-mismatch incident: alerts fire on dlq_depth > 1000. Mitigation: alert routes to the consumer team (NOT the dispatch team) so the right engineer triages; replay tool batches by 1000-message chunks with a kill switch; if the schema regression is in the producer (DISCO), rollback is the fix, NOT replay.

Stage by stage

The same 10 stages the workspace walks, answered.

01Clarifications

What would you ask before drawing a single box?

Typical clarifications to surface:

  • Pool / shared rides? Pool changes matching from 1:1 to 1:N optimisation — a different problem (the CAS predicate (driver_id, status='available') no longer works; status becomes (slots_used, slots_total)). Default: solo only.
  • Multi-region? Yes — region-pinned per city via the Global LB's GeoDNS pin table. NYC always served from us-east-1; SF from us-west-2. No cross-region dispatch on the request path.
  • Driver acceptance optional? Yes — drivers can decline. ~60% of offers a driver actually responds to are accepted, but ~half of offers expire unseen (backgrounded app, tunnel), so net conversion is ~31% per offer → ~3.2 offers issued per request on average.
  • Surge required? Yes — quoted in the offer envelope, must be stable for the offer's lifetime.
  • ETA accuracy? Within ~10% of actual, p95. Drives the choice of routing engine + DeepETA residual.
  • Latency target? rider tap → driver phone rings, p99 ≤ 3 s end-to-end.
  • What does "match" mean? A driver acknowledges the offer (ACCEPTED). EXPIRED / DECLINED / FAST_EXPIRED roll forward to the next candidate.

Assumptions to state:

  • 1.5 M concurrent online drivers globally; 4 s heartbeat → 375 K writes/s sustained, 1.875 M/s peak (5x).
  • 28 M trips/day → 324 req/s avg, ~1.6 K req/s peak globally; ~400 req/s NYC at Friday 11pm; event headroom 4 K req/s single-city.
  • ETA browse-to-request ratio ~10:1 → ~16 K ETA QPS peak (1.6 K peak request QPS x 10).
  • Offer TTL 15 s; ~3.2 offers issued per request on average; 75 K concurrent in-flight offers globally at peak.

02Functional reqs

What must this system actually do?

  • Rider taps "Request" with origin + destination + ride class. System returns either MATCHED with driver ETA + final price, or NO_DRIVER after retry budget exhausted.
  • Rider can GET /eta (browse) without committing — same envelope minus the dispatch.
  • Driver receives an offer push + WS message; POST /offers/:id/accept within the 15 s TTL transitions to ACCEPTED.
  • Driver location stream is continuous; the system must reflect a heartbeat in the dispatch index within 8 s (two cycles).
  • Surge multiplier is computed continuously and quoted inside the offer envelope; honored for the offer's lifetime regardless of refresh.
  • Offers stuck in OFFERED beyond their TTL are force-EXPIRED by the offer-reaper (out-of-band liveness).

03Non-functional

What must it promise about speed, uptime and correctness?

  • Availability: 99.95% on the dispatch path (rider request → offer issued). Dispatch is the product; downtime means revenue loss + driver churn.
  • Latency: p99 rider-tap → driver-phone-rings ≤ 3 s end-to-end. Within that: ETA p99 ≤ 200 ms; matcher p99 ≤ 700 ms; offer fan-out (push + WS) p99 ≤ 1.5 s.
  • Durability (offer state): in-region RPO ≤ 0 (Cassandra QUORUM RF=3 across 3 AZs survives single-AZ loss with no data lost), cross-region RPO ≤ 5 s (Kafka uReplicator), in-region RTO ≤ 30 s, cross-region RTO in minutes (we accept; cities are region-pinned).
  • Durability (driver location): RPO ≤ 5 s. A heartbeat refills the index in <2 cycles, so this is bounded by Kafka mirror lag, not by the KV.
  • Consistency: kRing candidate read can be eventual (replica-served); the offer-claim CAS on geoidx must be linearizable per driver record (no double-dispatch). The Cassandra QUORUM offer INSERT is sequential.
  • Fairness: per-driver offer cooldown after decline so the "best ETA" greedy doesn't starve one driver with offer spam.
  • Region pinning: every city is owned by exactly one region (pin table at the Global LB); a region failure drops that region's cities, not the global product.

04Capacity estimation

How much load and data does this have to hold?

Inputs: 1.5 M concurrent online drivers, 4 s heartbeat, 28 M trips/day, 5x peak multiplier, 3.2 offers/request, 15 s offer TTL, 10x ETA-browse multiplier.

Driver location ingest (the dominant write path).

  • Sustained: 1.5 M / 4 s = 375 K writes/s.
  • Peak: x 5 = 1.875 M writes/s. Per city (NYC, ~50 K concurrent): 12.5 K writes/s.
  • Per Redis shard: 1.875 M ÷ 64 shards × 2 ops/heartbeat (SET + ZADD) ≈ 58K ops/shard/s — feasible only with pipelined per-shard client (locingest.keyChoices). A non-pipelined client would saturate at ~50% of this.
  • Implication: this cannot go through a B-tree on disk. In-memory KV (Redis cluster, 64 shards, 2 read replicas/shard for hot-cell scale-out) plus a Kafka tee for durability and downstream consumers.

Rider requests.

  • Avg: 28 M / 86 400 = 324 req/s. Peak: x 5 = ~1.6 K req/s globally; ~400 req/s NYC; event headroom 4 K req/s single-city.
  • In-flight offers at any moment: 1.6 K x 15 s = ~24 K (x ~3 fanout = 75 K offer rows alive at peak globally).

Geo-index R:W asymmetry.

  • Reads: 1.6 K req/s × 1 kRing(k=2, 19 cells) = ~30 K cell scans/s, but coalesced (per-shard pipelining) to ~5 K Redis ops/s.
  • Writes: 1.875 M/s peak.
  • R:W = 1:125 — the load shape is the inverse of the typical web caching story. Bottleneck is write throughput, not cache hit rate.

ETA + surge cache.

  • ETA QPS peak: 1.6 K x 10 (browse multiplier) = ~16 K QPS. Surge cache reads = ETA QPS (one GET per envelope).
  • Surge cache cardinality: ~6 M H3 r7 cells globally; hot subset ~500 K.

Offer + ride storage (Cassandra).

  • Offer writes: ~5 K/s peak (1.6 K req x 3.2 offers). Idempotent INSERTs keyed on offer_id.
  • Ride state UPDATEs (OFFERED → ACCEPTED → IN_TRIP → COMPLETED): ~8 K/s peak. Combined ~13 K writes/s — well under Cassandra QUORUM ceilings (~50 K writes/s on a 256-vnode RF=3 cluster on commodity hardware).
  • Ride records: 28 M/day x 4 KB ≈ ~40 TB/yr; with TTL on offer-only rows (purged after 30 d) the live working set is ~1 yr of completed rides + ~75 K active offers.

Surge aggregator (Flink).

  • Input: 1.875 M heartbeat msg/s + 5 K dispatch event/s ≈ 1.88 M msg/s peak.
  • State: 6 M cells × 3 windows × ~3 KB/entry ≈ 54 GB working set.
  • Capacity: 16 TaskManagers × 4 GB RocksDB = 64 GB total — fits with 19% headroom; the 5x cardinality growth is the trigger to re-tier to hierarchical aggregation.
  • Output: 1 publish per cell per 10 s slide × 500 K hot cells = 50 K writes/s into the surge cache (batched).

End-to-end latency budget (p99 ≤ 3 s).

legbudget p99
rider TLS + GLB anycast route200 ms
GLB → region gateway50 ms
gateway auth + rate-limit80 ms
riderapi orchestrate60 ms
ETA (incl. surge + Maps)200 ms
matcher (kRing + score + CAS)250 ms
rides Cassandra QUORUM commit100 ms
outbox + push + WS delivery1500 ms (APNs tail dominates)
total~2.4 s p99, ~3 s p99.9

Cascading retry math (post-revision). Single-shot on the dispatch hot path; idempotency keys carry the retry budget at the rider layer.

  • e1 (rider→glb): 5000 ms total budget; rider exp-backoff retries with Idempotency-Key.
  • e3b (glb→gateway): 4000 ms Envoy overall-timeout deadline — the 1 retry shares this budget (not 2x it, unlike the per-attempt DISCO-internal edges below), so the whole hop still nests inside the rider's 5000 ms.
  • e4 (gateway→riderapi): 2500 ms Envoy overall-timeout deadline — the 1 retry shares this budget, so the hop nests inside glb's 4000 ms.
  • e6 (riderapi→eta): 500 ms, circuit-breaker (no retry doubles the budget); on open, riderapi proceeds with surge=1.0.
  • e9 (riderapi→disco): 1500 ms, single-shot (retryPolicy: "none"); rider's Idempotency-Key + riderapi's request_nonce LRU mean a retried client request joins the in-flight match instead of double-dispatching. Total: 1500 ms ≤ 2500 ms parent budget with 1000 ms slack for orchestration overhead.
  • e11a (disco→geoidx kRing read): 50 ms with fixed-retries → 100 ms worst-case.
  • e11b (disco→geoidx CAS-claim): 30 ms, single-shot; CAS itself fails closed on contention so a retry would risk a double-claim race.
  • e12 (disco→rides): 250 ms, fixed-retries → 500 ms worst-case (idempotent on offer_id).
  • e13 (disco→bus): 30 ms, fixed-retries → 60 ms worst-case; if both attempts fail, the outbox-relay re-drives from CDC.
  • DISCO inner total worst-case: 100 + 30 + 500 + 60 ≈ 690 ms — under the 1500 ms slot.

05API design

What does the outside world call, and what comes back?

GET /api/v1/eta?origin=40.75,-73.98&dest=40.78,-73.93
Authorization: Bearer <rider-jwt>

200 OK
{
  "eta_seconds": 540,
  "base_price_cents": 1820,
  "surge_multiplier": 1.4,
  "surge_version": "v42-1715000000",
  "ttl_seconds": 60
}
POST /api/v1/requests
Authorization: Bearer <rider-jwt>
Idempotency-Key: <client-uuid>

{
  "origin": {"lat": 40.75, "lng": -73.98},
  "dest":   {"lat": 40.78, "lng": -73.93},
  "ride_class": "uberx",
  "surge_version": "v42-1715000000"
}

202 Accepted
{
  "request_id": "req_01HZ...",
  "state": "SEARCHING"
}
# Matching is asynchronous (kRing scan → CAS-claim → offer), and a single request fans
# out ~3.2 offers because drivers decline/expire. The terminal result is PUSHED over the
# rider's existing WebSocket — it is NOT returned inline on this POST:
#
#   { "type": "offer", "request_id": "req_01HZ...", "offer": {
#       "offer_id": "ofr_01HZ...", "driver_eta_seconds": 220,
#       "final_price_cents": 2548, "surge_multiplier": 1.4,
#       "expires_at": "2026-05-07T15:00:30Z" } }
#   { "type": "no_driver", "request_id": "req_01HZ..." }   # search window exhausted
GET /api/v1/requests/:request_id            # poll fallback when the WS is down
Authorization: Bearer <rider-jwt>

200 OK
{ "request_id": "req_01HZ...", "state": "SEARCHING" | "MATCHED" | "NO_DRIVER" | "CANCELLED",
  "offer": { ... } }                        # offer present once state=MATCHED
POST /api/v1/offers/:offer_id/accept
Authorization: Bearer <driver-jwt>
Idempotency-Key: <driver-uuid>

200 OK
{ "state": "ACCEPTED", "ride_id": "rid_01HZ...", "rider_pickup_eta_s": 220 }

# Idempotent: a retry (tunnel reconnect) of an offer THIS driver already accepted
# replays 200 with the original ride_id, NOT 409 — the Idempotency-Key keys the
# claim so accept→drop-WS→retry can't race a fast-EXPIRE into a false conflict.
#
# Genuinely lost:
410 Gone     { "error": "offer_expired" }          # TTL / fast-EXPIRE reaped it
# Lost to another claimant or a stale ring epoch:
409 Conflict { "error": "offer_no_longer_valid" }  # claimed-by-other / stale-epoch
POST /api/v1/offers/:offer_id/decline
Authorization: Bearer <driver-jwt>

200 OK
{ "state": "DECLINED" }
# Idempotent: a re-declined offer replays 200. Decline immediately frees the
# geoidx driver claim and emits DECLINED to disco so riderapi rolls to the next
# candidate NOW — without waiting on the 15 s TTL reaper (the fast-rollforward
# economics depend on this; ~3.2 offers/request makes decline high-frequency).
409 Conflict { "error": "offer_not_offered" }      # already ACCEPTED/EXPIRED, not in OFFERED
POST /api/v1/requests/:request_id/cancel          # rider abandons a SEARCHING request
Authorization: Bearer <rider-jwt>

200 OK
{ "request_id": "req_01HZ...", "state": "CANCELLED" }
# Idempotent; tears down the in-flight match + any outstanding offer.
409 Conflict { "error": "already_matched" }        # a driver already ACCEPTED

Why version-stamp the surge? Without surge_version echoed on confirm, a rider who sees price X at offer time and price Y on retry can game the system; the gateway rejects mismatched versions and re-issues a fresh offer at the current price.

06Data model

What gets stored, and what is it looked up by?

offer (Cassandra, partition key offer_id):

fieldtypenotes
offer_iduuidPK; client-generated by DISCO
ride_iduuidFK to ride record once accepted
request_iduuidthe rider request this offer belongs to (trace correlation + cancel lookup)
rider_idbigint
driver_idbigint
stateenumOFFERED, ACCEPTED, DECLINED, EXPIRED, FAST_EXPIRED
surge_versiontextechoed back on confirm
epochbigintmonotonic ring-shard epoch — CAS rejects stale
expires_attimestampreaper sweeps
relayed_attimestamp?null until outbox-relay or primary publishes
h3_r9_originbigintfor analytics

request (Cassandra, partition key request_id): the durable rider-side request state riderapi persists so GET /requests/:id (poll) and POST /requests/:id/cancel resolve from any replica, not just the one holding the in-memory nonce LRU.

fieldtypenotes
request_iduuidPK; derived from the client Idempotency-Key
rider_idbigint
stateenumSEARCHING, MATCHED, NO_DRIVER, CANCELLED
origintextlat,lng
desttextlat,lng
ride_classtextuberx, ...
surge_versiontextquoted-envelope version
matched_offer_iduuid?null until state=MATCHED
created_attimestamp

ride (Cassandra, partition key ride_id): canonical post-acceptance record with status, pickup, drop, fare, timestamps.

offer_by_expiry (Cassandra, time-bucketed, TWCS): (time_bucket, offer_id) PK — the reaper scans this, NOT the main offer table, to avoid full-table tombstone walks.

geoidx (Redis): two keyspaces — driver:{id} → (h3_r9, lat, lng, status, ts, epoch) with TTL 30 s; cell:{h3_r9} → SortedSet<driver_id, score=ts> refreshed on each heartbeat.

surge (Redis): surge:{h3_r7} → (multiplier, version, computed_at), TTL 60 s.

Indexes & invariants:

  • (driver_id) UNIQUE enforced by single-key SET in geoidx (Redis cluster routes to one shard).
  • (offer_id) UNIQUE enforced by Cassandra partition-key write (idempotent INSERT).
  • (driver_id, status='available') is the CAS predicate at offer commit — fails closed if the driver was already claimed by another offer.

07High-level design

Which components handle a request, and in what order?

Architecture summary. Six lanes flow left-to-right:

  1. Region routing (clients → glb): rider/driver hits the closest anycast PoP; the GLB pin table routes to the correct city's home region.
  2. Edge (glb → gateway): Envoy at the regional edge does WAF, JWT, rate-limit (per-IP, per-user, per H3 r7 cell), and circuit-breaks slow upstreams.
  3. Sync request plane (gateway → riderapi → eta → disco → state): riderapi orchestrates the request, fanning to ETA (which fans to surge cache + Maps) for the offer envelope and DISCO for the actual match. DISCO is the geo-sharded matcher, sharded by H3 r8 supercell via embedded RingPop.
  4. State (geoidx + rides + surge): Redis cluster for the live driver index (write-dominant; durability lives in Kafka), Cassandra for canonical offer + ride records, Redis again for the surge multiplier read cache.
  5. Driver write plane (driver → glb → gateway → locingest → geoidx + bus): locingest is the heartbeat fast-path, isolated from the matcher so a slow DISCO cannot back-pressure ingest. Tees every sample to Kafka for durability + downstream surge.
  6. Async tier + recovery (bus → surgeagg → surge; bus → dlq for poison; rides → outboxrelay → bus for primary-publish-failure recovery; reaper → rides + bus for stuck-offer liveness): Flink windowed aggregator computes the per-cell surge multiplier on a 1-min sliding / 10-s sliding window and refreshes the surge cache. The outbox-relay, reaper, DLQ, and monitor sit alongside the hot path as the correctness + observability tier.

Read flow (rider tap → matched). rider → glb → gateway → riderapi (idempotent on Idempotency-Key) → riderapi fans out to eta (→ surge + maps in parallel) AND disco (→ in-process ring → geoidx kRing scan + CAS-claim → rides INSERT QUORUM → bus outbox + push fan-out) → POST returns 202 SEARCHING; the matched offer (or NO_DRIVER) is pushed to the rider over the existing WS.

Write flow (driver heartbeat). driver → glb → gateway (ws) → locingest → geoidx pipelined SET+ZADD (fire-and-forget) AND bus publish (async). Heartbeat dropped at any layer is replaced 4 s later — no retry, no amplification.

Driver accept flow. driver → glb → gateway → riderapi → disco (validates state + epoch) → rides CAS state OFFERED→ACCEPTED → bus publish ride.events.ACCEPTED → 200 to driver, push to rider via existing WS.

Decline / cancel flow. POST /offers/:id/decline (driver) and POST /requests/:id/cancel (rider) ride the same edges as accept and the original request respectively — driver decline reuses the driver → glb accept edge (e3-driver-glb-accept, labelled POST /offers/:id/{accept,decline}), rider cancel reuses e1 rider → glb. Both terminate at disco: decline frees the geoidx claim and emits DECLINED so riderapi rolls forward immediately rather than waiting on the reaper; cancel tears down the in-flight match.

Liveness (stuck offer). reaper (1Hz cron, leader-elected) → scans offer_by_expiry → CAS UPDATE OFFERED→EXPIRED → publishes EXPIRED to bus.dispatch.events. DISCO sees the EXPIRED, riderapi rolls forward to the next candidate.

Outbox recovery. outboxrelay tails Cassandra commit log via CDC → finds offers with relayed_at IS NULL after 5 s → republishes to bus.dispatch.events → stamps relayed_at. Idempotent on offer_id + state.

08Deep dives

Which part breaks first, and what do you do about it?

1. Why H3 r9 for matching, r7 for surge. Hexagons (vs squares/geohash) give equidistant neighbours so kRing(k) is an isotropic radius search — no directional bias. r9 cell ~150 m edge ~0.1 km^2 is the smallest unit at which "this driver and this rider are close enough" is meaningful in cities (one block-ish). r7 (~1.2 km edge, ~5 km^2) for surge because pricing wants stable smoothed regions, not per-block whiplash.

2. RingPop sharding by cell, not by city. City-sharding fails the moment one city's hot cell saturates one shard. r8 supercell sharding lets the load shed across cells; one MSG-cell stops one r8 supercell, not all of NYC. The fairness cost is gossip convergence latency during deploys (10 s typical), which we contain with the monotonic epoch on every offer.

3. Why the offer-claim is a CAS in geoidx, not a row-lock in Cassandra. Cassandra's lightweight transactions (LWT) cost 4x a regular write and pile up under load. The driver_id record in geoidx is single-shard, ms-level, and SET-CAS-able natively in Redis. We do CAS in Redis for "is this driver still available", then idempotent INSERT into Cassandra for the durable offer. The Cassandra write is the canonical commit; Redis is the scratchpad that prevents double-dispatch. The CAS-claim is part of the offer-commit saga (e11b) — its participant role is driver-claim, the failure mode is "claim-already-held", and the compensating action is "next candidate".

4. Why the outbox group spans Cassandra + Kafka + the relay worker. On dispatch, we want both the durable offer record (Cassandra) and the bus event (Kafka) to land. The cheapest correct pattern is the local-write-then-publish outbox: write Cassandra, publish Kafka with acks=all, mark relayed_at. If the Kafka publish fails, the outbox-relay worker tails Cassandra CDC, finds rows with NULL relayed_at after 5 s, and republishes idempotently. We rejected 2PC across Cassandra + Kafka — the coordinator is a SPOF in 2PC's PREPARE-COMMIT window, and one of the canonical retros for ride-hailing was a 2PC-coordinator stuck-locks incident. The relay worker is the load-bearing recovery path.

5. Hot cell at concert end (MSG, NYE). Pre-shard known-event cells into virtual sub-shards keyed on (h3_r9, hash(rider_id) mod N) ahead of the event window. Gateway-side admission limiter on the cell returns 429 Retry-After 5–30 s jittered before kRing fanout begins. KRing radius is capped (k≤3) so an unprovisioned cell does not exponential-fan-out. The geoidx hot-cell ZADD problem (one r9 cell's drivers all ZADD to one shard) is mitigated by the 2 read-replicas-per-shard story: kRing reads scale across replicas, leaving the leader free for ZADDs.

6. Stale surge during Kafka backlog. Surge multiplier is version-stamped; the offer envelope carries (multiplier, version); the gateway rejects confirms whose version is older than the offer (replay attack defense) AND honors the offered price for the offer TTL regardless of fresher refreshes (UX trust). The aggregator has a hard cap on per-cell multiplier delta per slide (≤ 0.5x) in both directions so a bad model can't globally triple prices AND a sudden supply unlock can't collapse a quoted offer to 1.0x mid-confirm.

7. Replica ghost-driver. KRing read can come from a Redis follower (eventual). The CAS-claim on driver record fails closed when another DISCO instance just claimed it; the matcher falls to the next candidate. The asymmetry is intentional: candidate-read is QPS-dominant and tolerates staleness; commit must be linearizable. Edge e11a (read) is fixed-retries; edge e11b (CAS) is single-shot.

8. Region failover. Cities are pinned to one region via the GLB pin table; cross-region replication is async via Kafka mirror at <5 s lag. A regional outage drops cities served from there. Riders in those cities see "no driver found" briefly until the GLB rotates the pin to the standby region (60s overlap window during the rotation lets both regions accept reads). In-flight offers ARE lost (RPO 5 s on offer state cross-region); the rider re-issues, idempotency-key dedupes against an empty match.

9. Driver-offline-mid-OFFER fast-EXPIRE. Naive 'cancel on disconnect' double-cancels legitimate transient blips (subway entry/exit). Production design: WS keepalive miss for >2 s on a driver holding an OFFERED offer triggers a fast-EXPIRE in DISCO — transition OFFERED → FAST_EXPIRED, free the driver claim in geoidx, riderapi rolls forward to the next candidate. Under 2 s the offer continues to ride out the legitimate WS hiccup. This compresses the "rider waits 15 s for next candidate" UX hit to <2 s.

09Trade-offs

What did this design cost, and what breaks at 10×?

What we accepted.

  • No global distributed transaction across geoidx + rides + bus. Outbox + monotonic epoch on the offer is what defends double-dispatch; we accept that under a hard partition we may emit two offers and reject the second at CAS, which costs one rider an extra second.
  • Surge can be stale by up to one window slide. The 60-s TTL on surge cache + 10-s slide means the worst case is ~70 s of stale multiplier on a regional Redis restart. We honor the quoted price on confirm regardless, so the rider is never charged a higher multiplier than they were quoted.
  • Async cross-region replication on offer state. RPO 5 s on in-flight offers means ~5 s of offers are lost on regional failover. The alternative — sync cross-region replication — adds 30–80 ms to every dispatch, which busts the 600 ms matcher budget. We chose RPO over latency.
  • Push-as-wakeup, WS-as-transport. APNs is best-effort; we accept that an APNs outage drops backgrounded iOS drivers temporarily; foregrounded drivers (the majority of the active pool) keep flowing through the WS.
  • kRing radius is capped (k ≤ 3). In a sparse rural area the matcher will return NO_DRIVER even if a driver exists 5 km away. The alternative — unbounded kRing expansion — DOSes the matcher in dense cells. Drivers in rural areas use a different matching policy (broadcast offer on the city scope) that is out of scope for this canonical.
  • DISCO call is single-shot (retryPolicy: 'none'). Trades retries for an idempotency-key narrative — riderapi's nonce LRU + rider-side Idempotency-Key are the budget-fitting alternative to caller-side retry. The gateway's 2500ms Envoy overall deadline spans its 1 retry and nests inside glb's 4000ms (itself inside the rider's 5000ms); DISCO's 1500ms single-shot inside that has 1000ms slack.

What breaks at 10x.

  • RingPop ring at 2 000 nodes has noticeably slower SWIM convergence; cell ownership churn at deploy time becomes a multi-minute event. Move to hierarchical rings (one ring per region with a top-level region router).
  • Single Kafka cluster on loc.heartbeats runs out of partition headroom past 1 200 partitions practically — re-tier to a dedicated cluster per region.
  • Cassandra QUORUM at 50 K offer/s starts needing per-region clusters with cross-region async; offer_id namespacing per region (uuid7 with region prefix) prevents collisions on the outbox topic.
  • Surge aggregator at 6 M+ keys crosses the RocksDB working-set boundary (54 GB) past 16 TM × 4 GB; switch to hierarchical aggregation (r9 → r7 roll-up at the worker, r5 → r3 at a downstream job).

Trade-offs we accepted (post-critic).

  • Per-driver fairness penalty ranker is in-process inside DISCO, not a separate node. The cooldown state is replicated through RingPop gossip on cell-ownership transfer. At the next scale tier (or if the fairness-vs-efficiency trade-off becomes a product KPI) this becomes its own feature-store-backed scorer.
  • Pool / shared rides are explicitly out of scope. The CAS predicate (driver_id, status='available') doesn't generalize; Pool requires status=(slots_used, slots_total) and a constrained-optimization matcher rather than a greedy. We do not retrofit; the canonical models solo-match.
  • The driver app receiving the push notification is modeled as a terminating arrow at External · Push. We do not draw a return arrow back to Driver Mobile because the actual delivery is a black box outside our system.
  • Single canonical store (Cassandra) for both offers and rides. At very large scale these split (offers in a TTL-aggressive store, rides in a long-retention OLAP-feeder). The split adds a join surface; not yet justified at this scale.
  • Cross-region RTO is in the minutes, not seconds. We accept this because cities are region-pinned via the GLB, so a region loss drops one region's cities (and only those cities), not the global product.

Primary sources

  • H3: Uber's Hexagonal Hierarchical Spatial Index (Uber Eng, 2018)
  • Scaling Uber's Real-time Market Platform — Matt Ranney, QCon 2015
  • Ringpop: Scalable, Fault-tolerant Application-Layer Sharding (Uber Eng, 2016)
  • Schemaless, Uber's Highly Available Datastore (Uber Eng)
  • DeepETA: How Uber Predicts Arrival Times (Uber Eng, 2022)
  • Streaming Lyft Ride Prices on Flink (Flink Forward SF, 2019)
  • Solving Dispatch in a Ridesharing Problem Space (Lyft Eng)
  • Real-time Data Infrastructure at Uber (arXiv:2104.00087)
  • Google SRE Workbook: Managing Cascading Failures

Now defend it

Reading a design is not the same as being able to hold one under questioning. The workspace asks the same questions an interviewer would, and the simulator disagrees with you when the diagram does not support the claim.

Work Uber / Lyft — Match Drivers and Riders yourself