Uber: Real-Time Dispatch, Ringpop and Schemaless
This page is one interview loop in three rounds, built on a real company's public engineering history. All three rounds design the same system: the part of Uber that matches riders with drivers in real time, and the part that stores every trip. Each round is an era. In each era, rapid growth forced Uber to rebuild one of those two parts while the business kept running.
| Round 1: Mid-level | Round 2: Senior | Round 3: Architect | |
|---|---|---|---|
| Era | About 2010–2014: a monolith and one PostgreSQL | About 2014–2018: dispatch on a gossip ring, trips on Schemaless | About 2018–today: hexagons, transactional stores and surviving a site loss |
| Level (Amazon) | SDE II (L5) | Senior SDE (L6) | Principal (L7) |
| Trips (published) | Growing about 20% a month in early 2014 | 1 billion trips in about six months (Dec 2015 to June 2016) | "40 million+ trips per day" across all products (Q4 2025) |
| Traffic we plan for | 20,000 drivers online at the peak: 5,000 location pings/s; 300,000 trips a day (assumptions) | 1M drivers online at the peak: 250,000 pings/s (assumption); ~5.5M trips a day (derived from published milestones) | 3M drivers and couriers online at the global peak: 750,000 pings/s (assumption) |
| Target | Keep matching and keep every trip while growth doubles the load every few months | No single database limit; add capacity without downtime | Survive losing a data center or region; consistency for multi-entity changes |
| Reading time | ~35 min | ~40 min | ~45 min |
You can start at any round. Rounds 2 and 3 open with a "Where we left off" summary that catches you up.
How to read a case study. Every claim about what Uber actually did comes from Uber's engineering blog, a talk by an Uber engineer, Uber's open-source code, or Uber's financial reports, and each round ends with a Sources list. We mark those claims Uber published or (published). Where Uber hasn't published the details, we say so and show a design that fits, labeled as ours. Numbers marked assumption are ours, chosen to make the arithmetic concrete. They are not Uber's internal figures.
Infrastructure note. For most of this history Uber ran its own data centers, first colocation space on the US West Coast and then more sites. By 2021 it was calling Google Cloud Spanner from its own data centers, and designed the links so its operational regions could be on-premises, AWS or GCP; in February 2023 it started moving from its own data centers to Oracle Cloud Infrastructure and Google Cloud (Uber, 2025). Where this page shows AWS services, that is a translation for this course, not Uber's setup.
Loop Opener: Why Uber?
Rebuilding the Plane While Flying It
Imagine a taxi company whose number of rides doubles every four months. Every piece of software you wrote last quarter is now running at twice the load it was designed for. You can't stop the business to rebuild: people are standing on street corners waiting for a car right now. That was Uber from about 2012 to 2016.
Two parts of the system felt the pressure most:
- Dispatch. The live marketplace. Drivers' phones report where they are every few seconds; riders ask for a car; dispatch picks a driver and sends an offer. Its data is fleeting: a driver's position from ten seconds ago is already useless.
- The trip store. The record of every trip: who, where, when, how much. Its data is permanent and valuable: billing, driver pay, support and fraud checks all read it. Losing a trip means losing money or trust.
Synthesizing vector architecture diagram...
Read it as three eras. Each rebuild answered a limit of the era before it: one database and one process (Round 1), then an eventually consistent, peer-to-peer design that was hard to reason about (Round 2).
A few words we'll use all page:
| Word | What it means on this page |
|---|---|
| Ping | One location report from a driver's phone: position, heading, speed and a timestamp. Uber's drivers sent one every 4 seconds (Uber, 2015). |
| Supply / demand | Uber's words for drivers (and couriers) and for riders' requests. |
| Shard | A slice of the data or the work, owned by one machine or one group of machines. |
| Consistent hash ring | A way to map keys to machines so that adding or removing a machine moves only a small share of keys. |
| Gossip | Machines telling a few random peers what they know, who tell a few more, until everyone knows. |
| Cell (Schemaless) | One immutable piece of data in Uber's Schemaless store: a JSON blob addressed by a row key, a column name and a version number. |
| Cell (geo grid) | One tile of a grid laid over the map (S2 squares, H3 hexagons). We always say which kind we mean. |
| Trigger | Schemaless's way of calling your code for every new cell, like a subscription to the store's change log. |
What Makes It Hard
- Growth outruns designs. Trips grew about 20% a month in early 2014 (Uber, 2015): that's 1.2^12 ≈ 8.9× in a year. A design with "a year of headroom" has about four months.
- Dispatch is live and heavy. Every online driver writes a position every 4 seconds, and every rider with the app open reads nearby cars. Reads and writes both scale with the fleet.
- Trip data can never be lost, and it is never finished. A trip gets its fare at the end, then billing attempts, then maybe a fare adjustment or a rating days later. Each arrives on its own schedule.
- Cities are wildly uneven. Uber's 2015 dispatch talk described its first, city-sharded design as getting harder to manage as cities were added, because some cities are big, some small, and only some have sharp load spikes.
The Question the Whole Loop Answers
How do we scale a live marketplace's dispatch and data layers step by step, without ever stopping the rides?
The answer grows every round:
- Round 1: split the one dispatch process by city, get the trip data off the shared database, and start breaking up the monolith.
- Round 2: shard dispatch by map cell over a self-organizing gossip ring (Ringpop), and move trips to an append-only store of immutable cells on sharded MySQL (Schemaless), with triggers that feed billing and analytics from the store's own log.
- Round 3: agree on one hexagonal grid (H3), give the stores transactions (Docstore, and Cloud Spanner for fulfillment), and plan for losing a whole site.
Round 1 · Mid-level · "Era 1: A Monolith and One Database"
~35 min · SDE II (L5) · about 2010–2014 · 20,000 drivers online at the peak and 300,000 trips a day (planning assumptions) · trips growing about 20% a month (published) · keep matching and keep every trip
R1.1 Establish Design Scope
The interviewer sets the scene: "It's early 2014. We started as a black-car service in San Francisco and now run in dozens of cities. There's a Python app with most of the business logic, a dispatch service that keeps the drivers in memory, and one PostgreSQL database that holds almost everything, trips included. Our trips grow about 20% a month. Keep it running, and tell us what to change first."
| We ask | Interviewer answers | What it changes in the design |
|---|---|---|
| What must the system do? | Request a ride, match a driver, run the trip, charge the rider. Uber (2015): the first codebase "solved our core business problems, which included connecting drivers with riders, billing, and payments". | Four flows; only matching is real-time (step 1.1). |
| Where are the drivers? | In memory, in the dispatch service. Uber described the old index as held in memory in a handful of processes (talk, 2015), and the old system "had to search through every single car" to find the nearby ones (blog, 2016). | Location data is fleeting, so memory is right. The scan is the problem (step 1.1). |
| How is dispatch split up? | By city: Uber's 2015 dispatch talk says the original dispatch system was sharded by city. | Good for isolation, uneven by nature (step 1.1). |
| Where are trips stored? | In the one PostgreSQL instance, with most other data. Trip data "was taking up the largest percentage, was (and still is) the fastest growing data, and also contributed to the most IOPS" (Uber, 2015). | The trip store must leave the shared database (step 1.2). |
| What's breaking? | Uber (2015): the trip store "was going to run out of steam both in terms of storage volume and IOPS by the end of the year, if not sooner". Adding a column or an index to the trips table "caused downtime". | A deadline measured in months, and schema changes we can't make (step 1.2). |
| What else hurts? | "Deploying the codebase meant deploying everything at once" (Uber, 2015). | Split the code along the flows (step 1.3). |
Out of scope for this round: sharding the trip store itself (Round 2), pooled rides, surge pricing, multi-region.
R1.2 Functional Requirements, Derived Step by Step
| Phrase from the problem | Operation |
|---|---|
| "Drivers report where they are" | sendLocation(driver, position) every 4 seconds while online |
| "Show me nearby cars" | nearbyDrivers(position, radius) for every rider with the app open |
| "Request a ride" | requestTrip(rider, pickup, product) returns a trip ID |
| "Match a driver" | Dispatch picks a candidate and sends an offer; the driver accepts or declines |
| "Run the trip" | The trip moves through its states: accepted, arriving, on trip, completed |
| "Charge the rider" | After the trip: compute the fare, charge the card, record the result |
Not yet: partitioned dispatch beyond one process per city, a sharded trip store, asynchronous downstream processing.
R1.3 Non-Functional Requirements: the Questions
We name each quality first; the numbers come in R1.7.
- Matching latency. How long from "request" to "a driver has an offer"? Riders are staring at a spinner.
- Freshness of positions. A position older than a couple of ping intervals is misleading.
- Write durability for trips. An acknowledged trip must survive a database crash; a completed trip that isn't stored can't be billed.
- Write availability. Uber (2016): "We needed write availability." When a write to PostgreSQL failed, the trip was parked in Redis to retry later.
- Growth per city. Each city grows at its own pace; one big city must not take down the others.
- Safe change. Adding a column, an index or a feature must not stop the business.
R1.4 The API
This API is illustrative: Uber's internal APIs of this era aren't public, so the paths and fields below are ours.
A driver's location ping (over a persistent connection from the driver app; shown as JSON):
json{ "type": "location", "driver_id": "d_2044", "lat": 37.77493, "lng": -122.41942, "heading_deg": 81, "speed_mps": 6.2, "accuracy_m": 8, "status": "AVAILABLE", "sent_at_ms": 1398350000000, "seq": 18221 }
seqis a per-driver counter the phone increments with every ping, so dispatch can drop a ping that arrives after a newer one.sent_at_msis the phone's clock. We use it to judge age, not to order pings (phones' clocks drift).
Request a ride
httpPOST /v1/trips HTTP/1.1 Authorization: Bearer <rider session token> Content-Type: application/json Idempotency-Key: 5d3c1f7e-2b8a-4c55-9e0a-7a1d2f9c3b61 { "product": "UBERX", "pickup": { "lat": 37.77493, "lng": -122.41942 }, "dropoff": { "lat": 37.78917, "lng": -122.40145 } }
httpHTTP/1.1 201 Created Content-Type: application/json { "trip_id": "t_7f3a91", "status": "REQUESTED", "poll_after_ms": 2000 }
Idempotency-Keylets the rider app retry after a timeout without creating a second trip. The server stores the key with the trip and returns the same trip for a repeat.
Status codes
| Code | Meaning |
|---|---|
201 Created | Trip created, dispatch started |
200 OK | A repeat of the same Idempotency-Key: the existing trip |
409 Conflict | The rider already has an active trip |
422 Unprocessable Entity | Pickup outside a city we serve |
503 Service Unavailable | Dispatch for this city is down; retry with backoff |
The trip's states
Synthesizing vector architecture diagram...
A trip moves forward only along these arrows. The loop between REQUESTED and OFFERED is dispatch trying the next driver. Everything after COMPLETED happens after the ride, on its own schedule.
Recap
- Two very different traffic types: frequent, tiny, disposable pings; and rare, valuable trip changes.
- The rider's request is idempotent; pings carry a sequence number.
R1.5 Design Evolution: From One Process and One Database to Split Responsibilities
Each step is a problem, your turn to think, the answer, and what it costs us.
Step 1.0: The Baseline
Synthesizing vector architecture diagram...
Two programs and one database. The dispatch process keeps every driver's position in memory and scans all of them for each request; the monolith writes every trip change into the one PostgreSQL.
Uber published this shape: "The early architecture of Uber consisted of a monolithic backend application written in Python that used Postgres for data persistence" (2016), and its 2015 dispatch talk describes dispatch as written almost entirely in Node.js.
Step 1.1: One Dispatch Process for Every City
The problem: one dispatch process holds every online driver and, for every "nearby cars" request and every ride request, scans all of them. It's single-threaded. Last month it pegged its CPU at the Friday evening peak, and when it crashed, every city lost dispatch at once. What would you do?
Loop: Design a Ride-Sharing Dispatch Service (Round 1: pings, offers, conditional accepts and timeouts for one city, in depth) · Loop: Design a Proximity Service (grid cells and radius search)
Step 1.2: Trip Writes Are Filling the One Database
The problem: the one PostgreSQL instance holds trips, users, payments and more. Trips are the biggest and fastest-growing table and cause the most disk I/O. Adding a column to trips now means downtime. At 20% growth a month, the disk fills before the end of the year.
What would you do?
Primitive: Database Sharding and Partition Keys · Loop primitive: Write-Ahead Log, fsync & Group Commit (Part 9 on the WAL in B-tree engines, and Part 10 on replication as a durability level)
Step 1.3: Every Deploy Is a Deploy of Everything
The problem: billing, payments, trips, user accounts and city configuration all live in one Python codebase. A bug in a new promo feature crashed the app servers on a Friday evening, and with them trip creation. Dozens of engineers are waiting on one release train. What would you do?
Primitive: Circuit Breaker, Bulkhead and Fault Tolerance
Round 1 Step Summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 1.0 | (baseline) | One dispatch process, one Python app, one PostgreSQL | One blast radius each |
| 1.1 | One dispatch process for every city | Dispatch sharded by city, a grid inside each city | Routing by city; uneven cities |
| 1.2 | Trip writes filling the database | A trip store behind one library, UUIDs, a Redis buffer, a plan to shard | A migration touching most code |
| 1.3 | Every deploy is a deploy of everything | Services along the flows | Remote calls that fail and retry |
R1.6 Architecture v1
Synthesizing vector architecture diagram...
Each city's drivers live in that city's dispatch process. All trip reads and writes pass through one library, so the store behind it can change. The Redis buffer keeps a trip when PostgreSQL rejects the write, at the price of that trip being invisible to billing until it's replayed.
Trace: a ride request in this era (our design around the published shape)
Synthesizing vector architecture diagram...
Every trip change is an in-place UPDATE of the same PostgreSQL row, which is why the trips table's write load grows faster than its row count. The post-trip work runs as one chain: if any step fails, the whole chain is retried.
Our AWS translation (not Uber's setup). Dispatch processes and the Python services on EC2, behind a load balancer; the one database as an RDS for PostgreSQL instance with a Multi-AZ standby in another AZ; the Redis buffer as a small ElastiCache cluster.
Sources for this round
- Schmidt, Project Mezzanine: The Great Migration, Uber blog, dated July 2015: ~20% monthly trip growth; storage volume and IOPS by the end of 2014; one PostgreSQL; trips the largest, fastest-growing and most IOPS; close to 100 services by early 2014; downtime for new columns and indexes; UUIDs; the
tripstoreswitch; mirrored writes and validation. - Thomsen, Designing Schemaless, Uber Engineering's Scalable Datastore Using MySQL, Uber blog, January 2016: "running out of database space" in early 2014; the five requirements; the Redis buffer; the synchronous post-trip processing.
- Klitzke, Why Uber Engineering Switched from Postgres to MySQL, Uber blog, July 2016: the Python monolith on Postgres; write amplification; verbose replication to the East Coast; the 9.2 replica corruption; hours-long upgrades.
- Haddad, Service-Oriented Architecture: Scaling the Uber Engineering Codebase As We Grow, Uber blog, September 2015: one offering in one city; deploying everything at once; 500+ services; "a distributed monolithic API".
- Ranney, Scaling Uber's Real-time Market Platform, QCon London, March 2015 (summarized by High Scalability): dispatch in Node.js; the original dispatch sharded by city; the old in-memory index in "a handful of processes"; pings every 4 seconds.
- Lozinski, How Ringpop from Uber Engineering Helps Distribute Your Application, Uber blog, February 2016: the earlier system that "had to search through every single car".
Everything else in this round (the grid inside each city, the API shapes, the service cut, the trace details and every number in R1.7 not marked published) is our design or our assumption.
R1.7 Numbers
Figures marked published come from the sources above; the rest are assumptions (January 2014 planning figures) or derived from them.
Dispatch load at the evening peak
| Quantity | Arithmetic | Result |
|---|---|---|
| Drivers online at the peak | Assumption | 20,000 |
| Location pings | 20,000 ÷ 4 s (the published interval) | 5,000/s |
| Riders with the app open | Assumption | 60,000 |
| "Nearby cars" requests | 60,000 ÷ 5 s (assumed refresh) | 12,000/s |
| Full scan, one process | 12,000 × 20,000 | 240M distance checks/s ≈ 24 cores at 10M checks/s per core (assumption) |
| Busiest city's process, full scan | 20% of both: 2,400/s × 4,000 drivers | 9.6M checks/s ≈ 1 core |
| Busiest city, with a 1 km² grid | 2,400/s × ~167 drivers in ~25 cells | ≈ 400,000 checks/s ≈ 4% of a core |
| Memory for all drivers | 20,000 × 200 B | 4 MB |
Trips
| Quantity | Arithmetic | Result |
|---|---|---|
| Trips a day | Assumption (January 2014) | 300,000 |
| Average trip requests | 300,000 ÷ 86,400 s | 3.5/s |
| Peak trip requests | × 3 (assumption) | 10.4/s |
| Row updates per trip | Assumption: states, driver, fare, billing attempts | 10 |
| Peak row updates | 10.4 × 10 | 104/s |
| Page changes per update (non-HOT) | 1 row version + 8 index entries (8 indexes, assumption) | 9 |
| Peak page changes | 104 × 9 | 936/s, each also in the WAL |
The runway (45 GB of trip data in January 2014: 300,000 trips × 30 days × 5 KB per trip including index entries and dead row versions, assumption; 1,200 GB free)
| Month of 2014 | Added that month (45 × 1.2^k GB) | Added so far |
|---|---|---|
| 1 (Feb) | 54 | 54 |
| 3 (Apr) | 78 | 197 |
| 6 (Jul) | 134 | 536 |
| 8 (Sep) | 194 | 891 |
| 9 (Oct) | 232 | 1,123 |
| 10 (Nov) | 279 | 1,402: past 1,200 |
The disk fills about 9.3 months in: early November 2014. By then the peak row updates have grown 1.2^9 ≈ 5.2×, to about 537/s, and page changes to about 4,830/s. Uber's trip store moved to Schemaless about a month before Halloween 2014 (published).
Synthesizing vector architecture diagram...
Compound growth hides the wall until late: more than half of the free space goes in the last three months before it fills.
R1.8 Trade-Offs
A monolith vs services, at this stage
| Keep the monolith | Split into services (Uber's choice) | |
|---|---|---|
| Deploys | One train; everything at once | Each team deploys its own service |
| Failures | A bug anywhere stops everything | A bug stops one flow |
| Calls | In-process, never time out | Remote: timeouts, retries, partial failure |
| Data | One schema, joins everywhere | Each service owns its data; no cross-service joins |
| When it wins | A few engineers, one city | Many teams, many cities, fast growth |
When to shard
| Signal | What it tells you |
|---|---|
| Disk runway under a year at the current growth rate | Start now: a migration takes months (Mezzanine's final push alone was 6 weeks, after months of preparation) |
| Schema changes need downtime | The table is too big for its database |
| Replicas lag or cost too much bandwidth | Replication is carrying more than reads need |
| One table dominates writes and I/O | Move that table first; leave the rest |
Synchronous standby or not? A synchronous PostgreSQL standby means a commit waits for the local WAL flush plus a round trip to the standby and its WAL write. With a standby on the other coast (an assumed 70 ms round trip), every trip update waits at least 70 ms longer. A standby in the same data center is cheap and protects against a lost disk, but not a lost site.
R1.9 Failure Modes
| Trigger | What you'd see | How the design responds |
|---|---|---|
| The database is at capacity | Slow queries, then timeouts on trip writes | Writes that fail go to the Redis buffer and are replayed later; billing waits for them (published). The durable fix is step 1.2's new store. |
| A dispatch process crashes | One city has no drivers on the map | Restart it. Driver positions come back within one ping interval (4 s), because every online driver pings again; nothing is lost that mattered. Trips in progress reload from the trip store. Other cities don't notice. |
| A city's process falls behind | Offers go out seconds late in that city | Its queue grows; shed "nearby cars" refreshes first (they're cosmetic) before ride requests. |
| A retried ride request | The rider tapped twice, or the app retried after a timeout | The Idempotency-Key returns the same trip. |
| A replica is corrupted | Wrong or duplicate rows on some replicas | Rebuild it from a fresh copy of the primary. Uber hit a PostgreSQL 9.2 bug after a promotion in which replicas misapplied WAL records, and fixed the replicas "by resyncing all of them from a new snapshot of the master" (published). |
R1.10 Pillar Check
| Pillar | What Round 1 covers |
|---|---|
| Reliability | Dispatch split by city so one crash stops one city; trip writes buffered when the database fails; idempotent ride requests REL 10 · REL 11 |
| Performance Efficiency | Positions in memory, a grid instead of a scan; the trip table's write amplification understood PERF 3 |
| Security | Each app talks only to the services it needs, with its own session token; driver positions are personal data and are kept only in memory SEC 7 |
| Operational Excellence | Services along the flows so teams deploy independently OPS 6 |
| Cost Optimization | Skipped this round: the constraint was time to the wall, not money. |
| Sustainability | Skipped this round: the fleet is small. |
R1.11 Round 1 Rubric and Follow-Ups
What a strong mid-level (L5) answer shows
- Separates fleeting location data (memory) from durable trip data (a database).
- Sizes the dispatch load and shows that CPU, not memory, is the limit, then fixes it with city sharding and a grid.
- Knows that replicas don't add write capacity.
- Computes a runway under compound growth and starts the migration in time.
- Puts all trip access behind one interface before changing the store.
- Splits the monolith along flows that fail differently.
Follow-up questions
-
"Why not keep trips in PostgreSQL and partition the table?" Answer: partitioning one table on one server splits the files, not the machine: disk, IOPS and the single writer stay the same. Uber's requirement was to "linearly add capacity by adding more servers", which needs shards on many servers.
-
"A trip starts in one city and ends in another. Which dispatch process owns it?" Answer: the pickup city's, for the whole trip. Ownership that moves mid-trip creates a window with two owners or none. The dropoff city only matters for pricing rules.
-
"The Redis buffer: what happens if Redis loses the trip too?" Answer: the trip is gone, which is why Uber wanted a store with write availability built in. Round 2's buffered writes put the backup copy in a second MySQL cluster, on disk.
Interview gotchas from this round's wrong answers
| Gotcha | Why it's wrong |
|---|---|
| "A bigger machine for dispatch" | The scan is CPU-bound and single-threaded; memory was never the limit. |
| "Read replicas for write load" | Every replica replays every write. |
| "A bigger database server" | Compound growth reaches the next wall in months. |
| "Deploy less often" | Bigger, riskier deploys, and every team still waits. |
| "Auto-increment trip IDs are fine" | They can't be generated on many shards; Uber had to replace them with UUIDs. |
Round 2 · Senior · "Era 2: Dispatch on a Gossip Ring, Trips on Schemaless"
~40 min · Senior SDE (L6) · about 2014–2018 · "many millions of trips per day" on six continents (published, 2016) · ~5.5M trips a day and 1M drivers online at the peak (planning figures) · 250,000 pings/s · no single database limit, capacity added without downtime
R2.0 Where We Left Off
This is what the candidate says aloud in the first 60 seconds of Round 2. If you're starting here, it's everything you need from Round 1.
Round 1 in 60 seconds. "Uber began as a Python monolith on one PostgreSQL, with a Node.js dispatch service that kept every online driver in memory. Drivers ping every 4 seconds. Scanning every driver for every request burned CPU in a single-threaded process, so dispatch was sharded by city, and inside a city drivers belong in a grid, not a list. Trips grew about 20% a month. Trip data was the biggest, fastest-growing and most I/O-heavy table in the shared database, and PostgreSQL's update model writes a new row version and new index entries for most updates, so the disk and IOPS would run out by late 2014. Replicas don't help writes. So all trip access went behind one
tripstorelibrary, trip IDs became UUIDs, failed writes were parked in Redis, and the monolith started splitting into services. Open costs: cities are uneven and each is still one process; the trip store still has to be sharded; and post-trip work runs as one synchronous chain."
Architecture v1, compact
Synthesizing vector architecture diagram...
Round 1 in one picture: dispatch split by city, trips behind one library, still one database underneath.
Round 1 step summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 1.1 | One dispatch process for every city | Sharded by city, a grid per city | Routing by city; uneven cities |
| 1.2 | Trip writes filling the database | tripstore, UUIDs, Redis buffer | A migration touching most code |
| 1.3 | Every deploy is a deploy of everything | Services along the flows | Remote calls that fail |
Open costs: big cities outgrow one process; the trip store is still one database; downstream work (billing, analytics) is synchronous and fragile.
R2.1 The Scope Raise
Interviewer: "It's late 2014 and growth hasn't slowed. We're in hundreds of cities on six continents. A city-to-process table that people edit by hand doesn't scale, and a big city still fits in one process only until next quarter. Dispatch must share work across many machines and survive losing some of them, without a central coordinator we'd have to keep alive. The trip database runs out of room in months. And billing, analytics and our other data centers need every change to every trip, without losing one."
| We ask | Interviewer answers | What it changes in the design |
|---|---|---|
| How many trips by 2016? | "Many millions of trips per day occurring across six different continents" (Uber, February 2016). Uber's CEO announced that the 2 billionth trip happened on June 18, 2016, about six months after the first billion at Christmas 2015. | About 1 billion in ~177 days ≈ 5.6M trips a day; we plan with 5.5M (R2.6). |
| How much location data? | Uber's 2015 dispatch talk: the new geospatial index had a design goal of a million writes a second, derived from drivers that update every 4 seconds, and runs on hundreds of processes. | We plan for 250,000 pings/s at our peak and keep the design goal as headroom (step 2.1). |
| How is dispatch organized now? | Split into a supply service (drivers and their state), a demand service (requests) and a matcher called DISCO, with a geospatial index of supply keyed by map cells (Uber talk, 2015). | Work is sharded by map cell, not by city (step 2.1). |
| What must the new trip store do? | Uber's five requirements (2016): add capacity by adding servers; write availability; "a way of notifying downstream dependencies"; secondary indexes; operational trust. | Steps 2.3 to 2.6. |
| Why not Cassandra, Riak or MongoDB? | Uber evaluated them and decided on "operational trust": whether "we would have the operational knowledge to immediately execute their fullest capabilities" (2016). | Build a thin layer on MySQL, which the team knew how to run (step 2.3). |
| How do downstream systems get changes today? | Many trip steps (billing, analytics) ran together, and "if any step failed, we had to retry all over again". The trip event system ran on Kafka 0.7 and "we were not able to run it lossless" (Uber, 2016). | A lossless change feed from the store itself (step 2.5). |
Scope change
| Round 1 | Round 2 | |
|---|---|---|
| Trips a day | ~300,000 (assumption) | ~5.5M (derived from published milestones) |
| Drivers online at the peak | 20,000 (assumption) | 1,000,000 (assumption) |
| Pings | 5,000/s | 250,000/s; design goal 1M writes/s (published) |
| Dispatch sharding | By city, fixed table | By map cell, on a self-organizing ring |
| Trip store | One PostgreSQL | Sharded, append-only, on MySQL |
| Downstream processing | One synchronous chain | Asynchronous, lossless, from the store's log |
The "Not yet" list from R1.2 comes back: partitioned dispatch and a sharded trip store are in scope now.
R2.2 What Breaks in the Round 1 Design
| Round 1 piece | What breaks at the new scale |
|---|---|
| A city-to-process table | Someone edits it by hand for every new city. A big city outgrows one process; a small one wastes one. A dead process takes its city down until someone moves it. |
| One PostgreSQL for trips | Disk and IOPS run out (Round 1's runway); schema changes need downtime; one writer. |
| Updates in place | Every trip change rewrites a row and its index entries, and two services updating the same trip row race each other. |
| A synchronous post-trip chain | One failed step (a card processor timeout) retries all steps, including those that succeeded. |
| A Redis buffer for failed writes | A buffered trip is invisible to billing, and it "did not scale" as Uber grew (published). |
R2.3 New Requirements and API Additions
The Schemaless API (Uber published these calls; the HTTP shapes below are ours). Uber described two kinds of access: random access to cells, and log-style access by shard.
| Call (published) | What it does |
|---|---|
put_cell(row_key, column_key, ref_key, cell) | Insert a cell. Rejected if that exact (row key, column, ref key) exists. |
get_cell(row_key, column_key, ref_key) | One exact version |
get_cell_latest(row_key, column_key) | The version with the highest ref key |
get_cells_for_shard(shard_no, location, limit) | Up to limit cells after location (an added ID or a timestamp) in one shard, plus the next location to ask for |
httpPUT /v1/instances/mezzanine/cells HTTP/1.1 Content-Type: application/json { "row_key": "8f1c2d9e-4b7a-4e21-9c1d-5a6b7c8d9e0f", "column": "STATUS", "ref_key": 2, "body": { "attempt": 2, "is_completed": true, "processor_ref": "ch_91ab" } }
httpHTTP/1.1 409 Conflict Content-Type: application/json { "error": "CELL_EXISTS", "row_key": "8f1c2d9e-4b7a-4e21-9c1d-5a6b7c8d9e0f", "column": "STATUS", "ref_key": 2 }
- A
409for an existing cell is how a writer learns it lost a race or is repeating itself. For a retry of the same write, the client treats it as success. - A write accepted while the shard's master is down returns success with a flag saying it isn't readable yet (Uber: "Schemaless can tell the client if the master is down"). We show it as
202 Accepted.
Trigger subscriptions (the concept is published; this config shape is ours). Uber: "Client programs link into the Schemaless triggers framework by configuring which Schemaless instance and which columns to poll data from."
yamltrigger_client: billing instance: mezzanine columns: [BASE] workers: 64 # at most one worker per shard; 4,096 shards offset_store: zookeeper # the added ID per shard; Uber: ZooKeeper or the instance itself poison_policy: max_attempts: 5 # then move the cell to a retry queue and continue stop_after_marked: 100 # too many marked cells means a systematic bug: stop and page
A secondary index (the concept is published: Uber's example is an index of a driver's trips, sharded on the driver's UUID, with the trip's city and creation time copied in; this definition is ours):
yamlindex: driver_trips column: BASE shard_field: driver_uuid # must be immutable and well spread: UUIDs are best fields: - driver_uuid - city_uuid # denormalized from the cell - trip_created_at # denormalized from the cell alias: driver_trips_v3 # a new version is backfilled, then the alias switches
Recap
- Writes are inserts of immutable cells; there is no update call.
- Reading "the current state" means reading the latest ref key.
- Downstream systems read the store's own ordered log, one shard at a time.
R2.4 Design Evolution: A Ring for Dispatch, Cells for Trips
Step 2.1: Which Machine Owns This Part of the Map?
The problem: dispatch now runs on hundreds of processes. A driver's ping must reach the process that holds that area's drivers; a rider's search must reach the processes that hold the areas around the pickup. Machines are added every week and die without warning. What would you do?
Primitive: Consistent Hashing · Primitive: Geospatial: Geohash, Quadtree and S2 · Drill: The Ring Rebalance That Crushed Node 07 (answered here: why adding 4 nodes doesn't shed 25% each without replica points, and why bounded loads are the wrong tool for stateful owners) · Loop: Design a Ride-Sharing Dispatch Service (Round 2: 500,000 location writes a second and batch matching)
Step 2.2: A Machine Dies, or Only Looks Dead
The problem: a dispatch process stops answering. Maybe it crashed; maybe it's in a one-second garbage-collection pause; maybe one network link is dropping packets. Every other process must agree on the ring, without a coordinator, and without evicting healthy machines every time one hiccups. What would you do?
Primitive: Gossip Protocol and Failure Detection · Drill: The GC Pause That Declared the Cluster Dead (answered here: why a fixed threshold flaps, and gossip vs a central coordinator)
Step 2.3: Trips Outgrow One Database
The problem: trips need a store that grows by adding servers, never needs downtime for a schema change, and lets many services add information to a trip at different times: the fare at the end, billing attempts, a fare adjustment days later, a rating. What would you do?
Primitive: Database Sharding and Partition Keys · Loop: Design a Distributed Key-Value Store (partitioning and replication from the inside)
How the move happened (published). Mezzanine ran in stages: trip IDs to UUIDs; a column layout for trips; a backfill from PostgreSQL; mirrored writes to both stores; rewriting every query; and "validation, validation, validation, and validation!", replaying queries against Schemaless and comparing, sampled to limit load on the busy PostgreSQL. The switch, about a month before Halloween 2014, "was a non-event". One rule makes a backfill and mirrored writes converge (our reading): both must write each trip version under the same ref key. Then whichever arrives second is rejected as a duplicate. A backfill that picked "latest + 1" instead would write an old version as the newest and hide the mirrored update, which is the append-only version of "backfills keep original write timestamps".
Step 2.4: A Master Dies Before It Replicates
The problem: a shard's master acknowledges a trip's BASE cell, and 50 ms later its disk dies. MySQL replicates asynchronously, so neither minion has the cell. The trip was never billed and now doesn't exist.
What would you do?
Loop primitive: Write-Ahead Log, fsync & Group Commit (Part 1: what "acknowledged" has to mean; Part 10: replication as a durability level) · Loop primitive: Idempotency & Effectively-Once Processing (Part 4: the late original; Part 5: idempotent by design)
Step 2.5: Billing and Analytics Need Every Change
The problem: when a trip's BASE cell is written, billing must charge the rider; analytics must load it; another data center must get a copy. Today one request does all of it in a row, and one failure retries everything. The old event system on Kafka 0.7 lost events.
What would you do?
Primitive: Change Data Capture and the Outbox Pattern · Loop primitive: Change Streams & the Transactional Outbox (Part 1: why writing twice fails; Part 5: reading the database's own log; Part 6: poison records and repair) · Loop primitive: Idempotency & Effectively-Once Processing (Part 6: queues and logs; Part 7: side effects you can't take back) · Drill: The Dual-Write That Broke Search Consistency (answered here: the dual write in the wrong answers, and polling vs log tailing above)
Step 2.6: Find a Driver's Trips
The problem: support needs "all trips of driver D in San Francisco last week". Trips are sharded by trip UUID across 4,096 shards. What would you do?
Primitive: Database Sharding and Partition Keys (global secondary indexes)
Step 2.7: Two Owners for One Trip
The problem: during a Tuesday deploy, dispatch nodes restart one after another. For a few seconds, two nodes both believe they own trip T's key. One records "driver A accepted"; the other, a moment later, records "offer expired, re-dispatch". The rider sees a car arriving, then the app says "finding your driver". What would you do?
Loop primitive: Leases, Fencing Tokens & Distributed Locks (Part 4: the pause, the zombie and the fencing token; Part 5: conditional writes are the check) · Primitive: Distributed Locks and Leases · Drill: The GC Pause That Corrupted Shared Storage (answered here: why a TTL lock doesn't survive a pause, and version checks vs an external lock) · Loop: Design a Ride-Sharing Dispatch Service (Round 1, step 1.3: the conditional accept written out in full)
Round 2 Step Summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 2.1 | Which machine owns this area? | Ringpop: consistent hash ring with 100 points per node, handle-or-forward; S2 level-12 cells as keys | A forwarding hop; fan-out on search |
| 2.2 | Dead or only slow? | SWIM: ping, ping-req, suspicion, incarnation numbers; flap damping | ~13.5 s to act; seconds of disagreement |
| 2.3 | Trips outgrow one database | Schemaless: immutable cells on 4,096 MySQL shards | Latest-version reads; no joins or transactions |
| 2.4 | A master dies before it replicates | Buffered writes on a random secondary; idempotent cells | 2 inserts and a delete per write |
| 2.5 | Every change, to every consumer | Triggers: per-shard log by added ID, at least once, poison queue | Trigger lag; idempotent handlers |
| 2.6 | Find a driver's trips | Sharded, denormalized, eventually consistent indexes | Stale reads; a repair path |
| 2.7 | Two owners for one trip | Version check at the store (LWT or unique cell) | Round trips; retries on hot entities |
R2.5 Architecture v2
Synthesizing vector architecture diagram...
Live dispatch state sits on Ringpop-sharded services backed by Cassandra; the permanent trip record goes to Schemaless when the trip ends, and everything downstream reads Schemaless's log through triggers.
Our AWS translation (not Uber's setup). Ringpop processes on EC2 in one Auto Scaling group per AZ, behind a Network Load Balancer for driver connections (priced in R2.6). Each Schemaless storage cluster as an RDS for MySQL instance with two read replicas in the other AZs: RDS read replicas replicate asynchronously, like Uber's minions, so buffered writes are still needed. The worker tier and trigger workers on EC2; offsets in a small replicated store.
Trace: a membership change during a search
Synthesizing vector architecture diagram...
Losing an owner loses nothing permanent: its cells refill from pings within one ping interval after the ring moves them. The cost is about 12 to 14 seconds of partial results for 0.3% of cells.
Sources for this round
- Thomsen, Designing Schemaless, January 2016: requirements; alternatives and "operational trust"; the data model; the trip columns; triggers; indexes (shard field, denormalization, "well below 20ms").
- Thomsen, The Architecture of Schemaless, January 2016: worker and storage nodes; 4,096 shards; one master and two minions across data centers; asynchronous replication; reads from the master by default; circuit breakers; buffered writes; idempotent writes; the entity table; MessagePack and zlib.
- Thomsen, Using Triggers On Schemaless, January 2016: the API; Schemaless as a partitioned log; added IDs identical on replicas; the triggers framework, leader, offsets, at-least-once, poison cells, scale.
- Schmidt, Project Mezzanine: The Great Migration, July 2015: 5 months to build; splitting servers in two; the migration stages; the switch in 2014.
- Nielsen and Johnsen, Code Migration in Production: Rewriting the Sharding Layer of Uber's Schemaless Datastore, February 2018: 40+ instances and many thousands of storage nodes in 2016; reads about 90% of traffic; the Go rewrite's results.
- Lozinski, How Ringpop from Uber Engineering Helps Distribute Your Application, February 2016: Ringpop's parts; Geospatial as its first use; many millions of trips a day on six continents.
- Uber, ringpop-go and ringpop-node source code (checked September 2026): 100 replica points, SWIM timeouts, suspect and faulty durations, flap-damping defaults, suspect members still reachable.
- Ranney, Scaling Uber's Real-time Market Platform, QCon London, March 2015: supply, demand and DISCO; S2 level 12 as the shard key; 1M writes/s design goal; hundreds of processes; Ringpop as AP; thousands of nodes; TChannel.
- Neerabail, Medisetty, Thangavelu and others, Uber's Fulfillment Platform: Ground-up Re-architecture, July 2021: the 2014 fulfillment stack (rt-demand, rt-supply, serial queues, Cassandra through a storage gateway, sagas) and its split-brain overwrites.
- Reuters, Uber reaches 2 billion rides six months after hitting its first billion, July 2016, reporting the CEO's announcement.
- MySQL 8.4 Reference Manual, Semisynchronous Replication; S2 Geometry, S2 Cell Statistics.
Uber hasn't published its trigger poll interval, its number of storage clusters, its peak factors, or how it handled index gaps, added-ID reuse after promotion, or fencing in the fulfillment stack; those parts of this round are ours.
R2.6 Numbers and Cost
All figures are assumptions unless marked published or derived.
Trips and the trip store
| Quantity | Arithmetic | Result |
|---|---|---|
| Trips a day | 1 billion in ~177 days (Christmas 2015 to June 18, 2016, published milestones) ≈ 5.6M; we plan with | 5.5M |
| Average trip rate | 5,500,000 ÷ 86,400 | 64/s |
| Peak trip rate | × 3 | ≈ 190/s |
| Cells per trip | BASE 1, STATUS 1.5, NOTES and adjustments 1.5 (assumption) | 4 |
| Cells a day | 5.5M × 4 | 22M |
| Peak cell writes | 22M ÷ 86,400 × 3 | ≈ 764/s |
| MySQL operations per cell (buffered) | 1 buffer insert + 1 primary insert + 1 buffer delete | 3 → ≈ 2,290/s at the peak |
| Index entries | 3 indexes on BASE × 5.5M a day | 16.5M a day, ≈ 573/s at the peak |
The trip store's write rate is modest. What forced sharding was volume, I/O on one machine, schema changes and write availability.
Storage
| Quantity | Arithmetic | Result |
|---|---|---|
| Stored per trip, compressed, with index entries | Assumption | 4.5 KB |
| Growth per copy | 5.5M × 4.5 KB | 24.75 GB a day |
| Growth, 3 copies (master + 2 minions) | × 3 | 74.25 GB a day ≈ 27.1 TB a year |
| Per shard, per copy | 24.75 GB ÷ 4,096 | ≈ 6 MB a day |
| Per storage cluster (32 clusters, 128 shards each) | 24.75 GB ÷ 32 | ≈ 0.77 GB a day ≈ 282 GB a year per master |
Doubling capacity means going from 32 to 64 clusters: each old cluster hands 64 of its 128 shards to a new one, "splitting each MySQL server in two".
Dispatch
| Quantity | Arithmetic | Result |
|---|---|---|
| Drivers online at the peak | Assumption | 1,000,000 |
| Pings | 1,000,000 ÷ 4 s | 250,000/s |
| Design goal (published) | 1M writes/s ÷ 0.25 per driver per s | room for 4M drivers: 4× headroom |
| Replica writes | 250,000 × 3 copies | 750,000/s |
| Searches at the peak | Assumption: open apps plus DISCO lookups | 60,000/s |
| Owner reads | 60,000 × 7 S2 level-12 cells touched by a 2 km circle (≈ 12.6 km² over cells averaging 5.07 km², plus the edges it crosses) | 420,000/s |
| Processes needed | 1,170,000 ops/s ÷ 10,000 per process (assumption) = 117; at 50% utilization, 234; so that 2 of 3 AZs carry it all, × 1.5 | 351: 117 per AZ |
Tail latency from pauses (Node.js garbage collection, assumptions). Suppose each process spends 0.5% of its time in 100 ms pauses. A request that lands in a pause waits anywhere from 0 to 100 ms, evenly spread. A search touches 8 processes (the one it lands on, plus 7 owners), so about 1 − 0.995⁸ ≈ 3.9% of searches meet a pause:
| Request | Share that meets a pause | p99 extra wait | p99.9 extra wait |
|---|---|---|---|
| One process (a ping, owner only) | 0.5% | 0 | 80 ms |
| A search over 8 processes | 3.9% | ≈ 75 ms | ≈ 97 ms |
(For the search's p99: the slowest 1% of all searches are the slowest quarter of the 3.9% that met a pause, which waited at least three-quarters of 100 ms.) Uber's answer, built into TChannel (published, 2015): backup requests with cross-server cancellation. Send the read to one replica, and if it hasn't answered after a short delay, send it to a second one; whichever answers first cancels the other.
The connection tier (AWS translation, us-east-1 list prices; 30-day month = 720 hours, used on every cost on this page). One Network Load Balancer with a TCP listener carries 1M persistent driver connections; TLS ends on our fleet.
| NLCU dimension (TCP) | Our load | NLCUs |
|---|---|---|
| New connections: 800/s each | 1M reconnecting once every 10 min (assumption) = 1,667/s | 2.1 |
| Active connections: 100,000 each | 1,000,000 | 10 |
| Processed bytes: 1 GB an hour each | pings 250,000 × 300 B = 75 MB/s, plus 25 MB/s to drivers = 100 MB/s = 360 GB an hour | 360 |
Billed on the largest dimension: 360 × $0.006 = $2.16 an hour × 720 = $1,555 a month, plus $0.0225 × 720 = $16 for the load balancer hour: ≈ $1,571 a month. (A TLS listener would count only 3,000 active connections per NLCU: 1M ÷ 3,000 = 334 NLCUs by connections, but the bill stays at 360 because bytes still dominate, and the TLS work moves to AWS.)
Cross-AZ transfer for pings (translation; $0.01/GB each way, so $0.02 for every GB that crosses)
| Stream | Arithmetic | A month |
|---|---|---|
| Forward to the owner (2/3 land in another AZ) | 250,000 × 300 B × 2/3 = 50 MB/s = 4,320 GB a day × 30 × $0.02 | $2,592 |
| Copy to 2 replicas in the other AZs | 250,000 × 300 B × 2 = 150 MB/s = 12,960 GB a day × 30 × $0.02 | $7,776 |
| Total | 200 MB/s | $10,368 |
The replication is what buys surviving an AZ loss without waiting 4 s for pings to refill. Two copies instead of three halves it; skipping replication entirely is also defensible for data that refills itself in 4 seconds.
Trigger polling (from step 2.5)
| Setup | Arithmetic | Queries/s |
|---|---|---|
| Poll every shard, every client, once a second | 4,096 × 20 clients | 81,920 (≈ 107 per written cell at the 764/s peak) |
| Batch per storage cluster | 32 × 20 | 640 |
| Average lag at 1 s polling | half the interval, plus processing | ≈ 0.5 s |
R2.7 Trade-Offs
Gossip vs a central coordinator for ownership
| Ringpop (gossip, chosen then) | Central shard manager (ZooKeeper, Helix) | |
|---|---|---|
| Single point of failure | None | The coordinator ensemble (a quorum) |
| Agreement after a change | Seconds; nodes briefly disagree | Immediate for everyone who can reach it |
| Scaling | Thousands of nodes (published), but gossip traffic and convergence grow with the ring | Coordinator load grows with changes, not requests |
| Failure detection | Built in (SWIM) | Sessions and heartbeats to the coordinator |
| Deploys | Every restart is a membership change | Assignments can be moved deliberately first |
| Uber's later view | Ring size capped pod size; long deploy cycles (2021) | Considered in 2021 as an incremental step |
Append-only cells vs update-in-place rows
| Update in place (Round 1) | Append-only cells (Schemaless) | |
|---|---|---|
| A change | Rewrites a row and its index entries | Inserts one cell at the end of the table |
| Two services changing one trip | Race to overwrite the row | Different columns, or a rejected duplicate version |
| History | Gone unless you add an audit table | Kept by design |
| Change feed | A separate system (dual write) | The table itself, in added_id order |
| Reading current state | One row | The latest ref key per column |
| Storage | Current version only | Every version, forever (with rare deletes) |
Build our own store vs adopt one. Uber's reason to build was "operational trust": a 3 am page on the trip store must be something the team can fix, and it knew MySQL. The thin layer took about 5 months. The long-term costs were an in-house system to staff and a "restrictive API" (Uber, 2021). Uber also noted, in 2016, that it had since "adopted both Cassandra and Riak with success in other areas". The honest rule: build when no product meets a hard requirement and you can staff it for years.
Triggers vs a Kafka stream (a translation question). Kafka (Amazon MSK on AWS) gives consumer groups, retention and replay for free, but it's a second copy of every change and a dual write unless it's fed from the database log. If we did use MSK, we'd set unclean.leader.election.enable to false (MSK's default is to allow unclean leader election, which can lose acknowledged messages, unless tiered storage is on, which requires it to be false), produce with acks=all and set min.insync.replicas to 2 with 3 replicas.
R2.8 Failure Modes
| Trigger | What you'd see | How the design responds |
|---|---|---|
| Gossip flap storm | Many nodes flipping suspect and alive; cells moving back and forth; membership checksums disagree | Suspicion and refutation absorb pauses under ~11.5 s; flap damping suppresses repeat offenders; alarm on faulty transitions per minute and on the number of distinct checksums. Look for the common cause (a bad network switch, a GC-heavy deploy) before touching timeouts. |
| Trigger lag | Billing or warehouse loads minutes behind | Per-shard lag metric (now minus created_at of the last processed cell); add workers (up to one per shard); check for a poison cell holding a shard; batch polls per cluster. |
| A poison batch | A shard stops advancing; the same cells fail over and over | Mark and move failing cells to the retry queue after N attempts (published). If a batch fails and we can't tell which cell, split it in halves until the bad cell is isolated, move it aside, and let the rest through. Stop and page when marked cells pass the threshold. |
| A shard master dies | Writes for its shards go to buffer tables; reads fall back to minions and may be stale | Buffered writes keep writes available (published); circuit breakers send reads to minions (published). Promote a minion, bump its auto-increment above the highest stored trigger offset plus a margin (keeping a sentinel row there on MySQL 5.7), or rewind that shard's trigger offsets if you can't bump, and let the buffer tables replay. Never let the old master take writes again: revoke its application credentials or set it read-only before the promoted minion takes traffic, so there's one writer per shard. |
| Split brain in dispatch | Two owners during a deploy or failover; conflicting trip states | Version checks at the store (step 2.7); roll deploys slowly enough that the ring converges between batches. |
| Ghost drivers | A driver in a tunnel still shows as available; offers time out | A driver is a candidate only if its last ping is under 10 s old (2.5 intervals, our choice); an offer the app doesn't acknowledge within a few seconds goes to the next driver. |
| An idempotency check reads a lagging minion | A trip is charged twice | Idempotency reads always go to the master; the card processor gets an idempotency key per (trip, attempt). |
| A delayed original write outlives its retry | The client library retried on another worker; the first attempt arrives later | Harmless: same (row, column, ref key), so the second arrival is rejected. |
R2.9 Production Gotchas
| Gotcha | Why it hurts | What we do |
|---|---|---|
In-place UPDATEs on sharded tables | Two services race on one row; each update rewrites index entries; no history | Append a new version; split columns by who writes them |
| Synchronous secondary indexes on the write path | Cross-shard 2PC on every write | Sharded indexes updated after the cell, with a repair path |
| Central locks for geo partitioning | A lock service in the path of 250,000 pings a second | Ownership by consistent hashing; fence durable writes with versions |
| Treating a buffered write as a win | Two writers of the same version both hear "stored" while the master is down | Exclusive decisions wait until the cell has replicated to a primary minion |
| Idempotency checks on replicas | A lagging replica says "not done yet" | Read the master for any "did I already do this?" check, and a check can still miss a buffered cell, so the card processor's key is the real guard |
| Trigger offsets across a master promotion | Reused added IDs are skipped | Bump the counter past the highest stored offset (sentinel row before MySQL 8.0), or rewind offsets if you can't; idempotent handlers |
| Fast failure-detection timeouts | Pauses become evictions and flap storms | Suspicion with refutation; damping; alarms on churn |
R2.10 Pillar Check
| Pillar | What Round 2 adds |
|---|---|
| Reliability | No coordinator to lose; SWIM with suspicion; buffered writes on a second cluster; minions in other data centers; idempotent writes and retries; the trip store migrated with mirrored writes and validation REL 11 · REL 4 · REL 8 |
| Performance Efficiency | Consistent hashing with 100 points per node; S2 cells as keys; added_id keeps inserts sequential; backup requests cut the pause tail PERF 3 · PERF 2 |
| Security | Trips are classified as sensitive personal data (they reveal where people go), kept in one store with its history; pings stay in memory and never enter the trip store SEC 7; each service reaches Schemaless with its own credentials, allowed to write only the columns it owns, which the column split makes natural (our design) SEC 3 |
| Cost Optimization | Adding capacity by moving shards instead of buying bigger machines; the price of cross-AZ ping replication, known and chosen COST 8 · COST 6 |
| Operational Excellence | Signals for ring churn, trigger lag per shard, buffer-table depth and marked poison cells; a written promotion procedure OPS 8 · OPS 10 |
| Sustainability | The Go rewrite of the Schemaless workers cut CPU by more than 85% for the same traffic (published) SUS 3 |
R2.11 Round 2 Rubric and Follow-Ups
What a senior (L6) answer adds over L5
- Explains consistent hashing with replica points, and why neither vnodes nor bounded loads fix one hot key.
- Walks through SWIM's ping, ping-req, suspicion and incarnation numbers, with a timing chain.
- Designs an append-only cell store and says what it costs (latest-version reads, growth, no transactions).
- Knows asynchronous replication loses acknowledged writes, and how buffered writes and idempotent cells close that gap.
- Builds a change feed from the store's own log, with at-least-once handlers and a poison-cell path; compares polling with log tailing.
- Fences split-brain writes at the store with versions, not in memory.
Follow-up questions
-
"Why did Uber shard dispatch by S2 cell instead of by city?" Answer: city shards are as uneven as cities, and a big city is one shard. Cells are small and many, so hashing spreads them evenly across hundreds of processes, and a big city simply uses more processes. The price is fan-out: a search touches several cells' owners.
-
"A trigger function charges a card, then crashes before recording the result. What happens?" Answer: the cell is delivered again (at least once). The function reads the latest
STATUSfrom the master, sees no result, and charges again, unless the charge carries an idempotency key the processor recognizes. That key is what makes a side effect you can't take back safe to retry. -
"Could Schemaless use semisynchronous replication instead of buffered writes?" Answer: it would protect acknowledged writes while a replica keeps up, but it falls back to asynchronous after the timeout (10 s by default), and it gives no write availability while the master is down. Buffered writes give both, at twice the inserts.
Interview gotchas from this round's wrong answers
| Gotcha | Why it's wrong |
|---|---|
| "A central coordinator routes every request" | A bottleneck and a single point of failure on the hottest path. |
| "Hash modulo N" | Changing N moves almost every key. |
| "Three missed heartbeats means dead" | Pauses and lossy links turn into flap storms. |
| "Semisync means no loss" | It switches itself off after a timeout. |
| "Publish an event after the write" | A crash in between loses the event: a dual write. |
| "Ringpop guarantees one owner" | Only after membership converges; fence at the store. |
Round 3 · Architect · "Era 3: Hexagons, a New Datastore and Surviving a Site Loss"
~45 min · Principal (L7) · about 2018–today · "40 million+ trips per day" and 9.7M monthly drivers and couriers (published, Q4 2025) · 3M online at the global peak and 750,000 pings/s (planning figures) · one grid, transactions, survive losing a site
R3.0 Where We Left Off
Round 2 in 60 seconds. "Dispatch was sharded by S2 map cell on Ringpop: every process embeds a SWIM membership protocol and a consistent hash ring with 100 points per node, and any process can take a request and forward it to the key's owner. Failure detection uses ping, ping-req through 3 others and a 5-second suspicion window with incarnation numbers, so pauses don't evict nodes, but a real crash takes about 13 seconds to act on and every deploy is a membership change. Live trip and driver state sat in Ringpop-owned services with serial queues, stored in Cassandra, with sagas for multi-entity changes; split brain during deploys and failovers let last-write-wins overwrite state. Trips moved to Schemaless: immutable JSON cells keyed by row, column and version on 4,096 MySQL shards, with buffered writes on a random second cluster so a master failure loses nothing acknowledged, triggers that read each shard's log in added-ID order at least once, and sharded, eventually consistent indexes. Open costs: the ring's size caps how big a city can get; there are no transactions across entities; every team draws its own map grid; and a whole site can still go down."
Architecture v2, compact
Synthesizing vector architecture diagram...
Round 2 in one picture: an AP ring for live state, an append-only store for the record, and a log-driven pipeline downstream.
Round 2 step summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 2.1 | Who owns this area? | Ringpop ring, S2 cells | A hop; fan-out |
| 2.2 | Dead or slow? | SWIM with suspicion | ~13.5 s to act |
| 2.3 | Trips outgrow one database | Schemaless cells on MySQL | No transactions |
| 2.4 | Master dies before replicating | Buffered writes | 2 inserts and a delete |
| 2.5 | Every change downstream | Triggers | Lag; idempotent handlers |
| 2.6 | Trips by driver | Sharded async indexes | Staleness |
| 2.7 | Two owners | Version checks | Retries |
Open costs: ring size limits, split-brain overwrites, sagas instead of transactions, many map grids, and one site's failure still stranding its cities.
R3.1 The Scope Raise
Interviewer: "It's 2019. We're not just rides any more: food delivery with a restaurant in the middle, freight, reservations, airport queues, a bus driver with a dozen passengers. Every product team has drawn its own map zones. Our fulfillment stack is an AP ring with sagas, and debugging it is a nightmare. Our trip store can't do transactions. And we run on our own data centers and two failure domains; losing one must not stop the marketplace."
| We ask | Interviewer answers | What it changes in the design |
|---|---|---|
| How big is fulfillment? | Uber (2021): "more than a million concurrent users and billions of trips per year across over ten thousand cities", "billions of database transactions a day", with 500+ developers building "120+ unique fulfillment flows". Today: more than 40 million trips a day (Q4 2025). | ≈ 472 trips/s on average (R3.6), and a platform many teams extend. |
| What's wrong with the map zones? | Uber (2018): zones drawn by hand or by postal code have "unusual shapes and sizes" and "require frequent updating"; it needed a grid with "comparable shapes and sizes across the cities". | One shared grid (step 3.1). |
| What's wrong with Schemaless? | Uber (2021): "the restrictive API and modeling capabilities made it hard for users to use as a general-purpose database". Teams turned to Cassandra, whose "eventual consistency ... ended up impeding developer productivity". | A transactional store built from Schemaless (step 3.2). |
| What must fulfillment guarantee? | Uber's 2021 requirements: "Strong consistency for single-row, multi-row, and multi-row multi-table transactions across regions"; region, zone or intermittent failures "should not cause data loss"; "at least 99.99% adherence to SLA". | Transactions for live state, replicated across regions synchronously (step 3.3). |
| How are sites laid out? | "Most of Uber's architecture ran on 2 failure domains with 2 independent regions", with cities grouped into pods inside each (Uber, 2021). | A plan for losing one of two (step 3.4). |
| Own hardware or cloud? | Uber ran its own data centers, reached Google's Cloud Spanner over dedicated interconnects from 2021, and in February 2023 began "migrating from on-premise data centers to the cloud with Oracle Cloud Infrastructure and Google Cloud Platform" (Uber, 2025). | What moves, and what that changes (step 3.5). |
Scope change
| Round 2 | Round 3 | |
|---|---|---|
| Products | Rides | Rides, delivery, freight and more: 120+ flows |
| Trips | ~5.5M a day | 40M+ a day (all products, published) |
| Map grid | S2 cells in dispatch, hand-drawn zones elsewhere | One hexagonal grid, H3 |
| Live state | AP ring, Cassandra, sagas | Transactions, strongly consistent across regions |
| Record store | Schemaless (no transactions) | Docstore: transactions per partition, Raft |
| Site loss | Two failure domains, split brain possible | Survive a region without losing acknowledged state |
R3.2 What Breaks in the Round 2 Design
| Round 2 piece | What breaks |
|---|---|
| Many map grids | Surge in one team's zones and supply in another's can't be compared; zones drift as cities change. |
| An AP ring for live state | Split brain during deploys and failovers overwrites state; pods can't grow past the ring's size; Uber measured the old stack at "only 20 online drivers/couriers per core" (2021). |
| Sagas across entities | Between steps the system is "in an internally inconsistent state", and failed compensations "often required manual intervention" (Uber, 2021). |
| A store without transactions | Every team reinvents read-check-append with versions, or moves to a store with weaker guarantees. |
| Two failure domains with asynchronous copies | A region loss can lose the last seconds of live state, and the survivor must carry everything. |
R3.3 New Requirements and API Additions
A shared grid library (H3; function names from the open-source library).
| Call | What it does |
|---|---|
latLngToCell(lat, lng, res) | The 64-bit H3 index of the hexagon containing a point, at resolution 0 (coarsest) to 15 (finest); geoToH3 before H3 v4 |
cellToLatLng(cell), cellToBoundary(cell) | The center and outline of a cell |
gridDisk(cell, k) | The cell and every cell within k steps: 1 + 3k(k+1) cells (kRing before v4) |
cellToParent(cell, res) | The coarser cell, by truncating the index |
One transaction for "driver accepts an offer" (a design in the style of Uber's fulfillment platform; SQL shown for clarity, not Uber's schema).
sqlBEGIN; SELECT owner_epoch FROM city_owners WHERE city_id = 'san_francisco' FOR UPDATE; -- must be 41, or abort -- the offer must still be open and meant for this driver UPDATE offers SET state = 'ACCEPTED' WHERE offer_id = 'o_55' AND state = 'OPEN' AND driver_id = 'd_2044'; -- each UPDATE must touch exactly 1 row, or abort UPDATE trips SET state = 'ACCEPTED', driver_id = 'd_2044', version = version + 1 WHERE trip_id = 't_7f3a91' AND state = 'OFFERED'; UPDATE supplies SET active_jobs = active_jobs + 1, version = version + 1 WHERE driver_id = 'd_2044' AND active_jobs < max_jobs; -- post-commit work, committed atomically with the state change INSERT INTO late_tasks (shard_id, task_id, kind, payload) VALUES (17, 'k_91', 'NOTIFY_RIDER', '{"trip_id":"t_7f3a91"}'); COMMIT;
All three entities change together or not at all, and the notification is recorded in the same commit, so it can't be lost or sent for a change that didn't happen.
Site ownership controls (our design).
yamlcity: san_francisco home_region: region_a standby_region: region_b owner_epoch: 41 # bumped on every move; every write for this city carries it failover: detection: outside_probes # health checks run from outside both regions approval: on_call # a human confirms before cities move data_residency: us_only # the standby must satisfy the same rules
R3.4 Design Evolution: One Grid, Transactions, and a Plan for Losing a Site
Step 3.1: Every Team Has Its Own Map Zones
The problem: pricing measures supply and demand in hand-drawn zones; dispatch uses S2 cells; analytics uses postal codes. A zone boundary moved last month and a year of surge data no longer lines up. Nobody can compare "demand here" with "supply here". What would you do?
Primitive: Geospatial: Geohash, Quadtree and S2 · Loop: Design a Proximity Service (Round 2: cell sizes, dense and sparse areas) · Loop: Design a Ride-Sharing Dispatch Service (Round 2, step 2.5: surge by zone)
Step 3.2: The Record Store Needs Transactions
The problem: new products need to change several records together (an order, its items, the courier's plan) and read their own writes. Schemaless offers immutable cells and no transactions; teams that moved to Cassandra struggle with eventual consistency. What would you do?
Primitive: Distributed Consensus: Raft and Paxos · Loop primitive: LSM-Trees & Compaction (MyRocks is an LSM engine under MySQL; Part 7 compares B-trees and LSM trees) · Drill: The Network Partition That Elected Two Leaders (answered here: what a 2-node minority can commit, and 3 nodes vs 5)
Step 3.3: Live Dispatch State Needs Transactions Across Regions
The problem: a driver accepts a batch of three delivery offers. Three orders and the driver's plan must change together. On the Ringpop stack this is a saga: propose on each entity, then commit, or run compensations. A failure mid-way leaves the driver with two of three orders, and someone fixes it by hand. What would you do?
Primitive: Two-Phase Commit and Saga Orchestration · Primitive: Database Isolation Levels, ACID and Concurrency Anomalies · Loop primitive: Change Streams & the Transactional Outbox (Part 2: the outbox in one transaction; Part 3: the relay and at-least-once delivery) · Loop: Transactional Outbox and Event-Driven Ledger
Step 3.4: A Whole Site Goes Down Mid-Trip
The problem: at 18:00 on a Friday, region A stops answering. Its cities have hundreds of thousands of trips in progress. Region B is healthy but running its own evening peak. What would you do?
Primitive: Cloud Disaster Recovery and Multi-Region Active-Active · Loop primitive: Leases, Fencing Tokens & Distributed Locks (Part 7: leader election and single writers) · Loop: Design a Ride-Sharing Dispatch Service (Round 3, step 3.6: phones resyncing trips after a region loss, in detail) · Drill: The Booking That Existed in Frankfurt but Not in Virginia (answered here: the double assignment from replication lag, and why not everything synchronous)
Step 3.5: Our Own Data Centers or the Cloud?
The problem: the company owns data centers with years of tooling around them. Some systems already call a cloud database over dedicated links. Leadership asks: move everything to the cloud, stay, or mix? What would you do?
Primitive: Cloud Disaster Recovery and Multi-Region Active-Active
Round 3 Step Summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 3.1 | Every team has its own zones | H3: hexagonal, hierarchical, 16 resolutions | Migration; approximate rollups |
| 3.2 | The record store needs transactions | Docstore: Raft partitions, strict serializability per partition, Flux CDC; MyRocks | Partition-key modeling |
| 3.3 | Live state needs transactions across regions | Fulfillment on Spanner multi-region; statecharts; LATE outbox | Cost 3–4× single region; cross-region commits |
| 3.4 | A site goes down mid-trip | Synchronous live state, outside detection, gated move, epochs, phone digests | Standby capacity; minutes to fail over |
| 3.5 | Own data centers or cloud | Hybrid, then OCI and Google Cloud (2023) | A migration; a transfer bill |
R3.5 Global Architecture
Synthesizing vector architecture diagram...
Each region runs everything for the cities it owns. Live, contended state commits synchronously across regions; the record store replicates Raft groups within a region and copies across regions; failover decisions and city ownership live outside both.
Our AWS translation (not Uber's setup). Two US Regions, us-east-1 and us-west-2, three AZs each. Dispatch and marketplace services on EC2 in Auto Scaling groups per AZ, behind Network Load Balancers. Live fulfillment state in an Aurora DSQL multi-Region cluster (the two Regions plus a witness Region). The record store as MySQL: RDS for MySQL Multi-AZ DB clusters (a writer and two readable standbys in three AZs, replicating semisynchronously), with our own sharding layer. City ownership in Route 53 Application Recovery Controller routing controls, with Route 53 health checks probing from outside. CloudWatch for alarms.
Trace: a driver accepts an offer
Synthesizing vector architecture diagram...
The accept, the trip, the driver's plan and the notification task commit as one. The notification may be sent twice (at least once), never zero times.
Trace: region A fails
Synthesizing vector architecture diagram...
The store already had every acknowledged change, so the digest mostly confirms what B knows. The epoch keeps a recovering region A from writing over B.
Sources for this round
- Brodsky, H3: Uber's Hexagonal Hierarchical Spatial Index, Uber blog, June 2018: why a grid; surge in hexagons; one neighbor distance; icosahedron, 122 base cells, 12 pentagons; 16 resolutions, one-seventh area; approximate containment.
- H3 documentation, Tables of Cell Statistics Across Resolutions (checked September 2026): average areas and edge lengths.
- Chatterjee, Chaudhary and Tariq, Evolving Schemaless into a Distributed SQL Database, Uber blog, February 2021: Docstore's motivation, layers, partitions of 3–5 nodes, Raft, strict serializability per partition, CP, MySQL transactions.
- Uber, MySQL to MyRocks Migration in Uber's Distributed Datastores, Uber blog, September 2022: shared topology, Raft within a region, replication across regions, tens of petabytes and tens of millions of requests a second, MyRocks since 2019.
- Pozniansky and others, How Uber Serves Over 40 Million Reads Per Second from Online Storage Using an Integrated Cache, Uber blog, February 2024: one leader and two followers; Flux CDC.
- Neerabail, Medisetty, Thangavelu and others, Uber's Fulfillment Platform: Ground-up Re-architecture, Uber blog, July 2021: scale; the old stack's problems; requirements; options considered; Spanner
nam3; LATE. - He, Medisetty and others, Building Uber's Fulfillment Platform for Planet-Scale using Google Cloud Spanner, Uber blog, September 2021: 20 drivers per core; evaluated databases; data model; LATE tailers; hybrid networking; cost split; snapshot then delete.
- Uber, Adopting Arm at Scale: Bootstrapping Infrastructure, Uber blog, February 2025: the February 2023 move to Oracle Cloud Infrastructure and Google Cloud.
- Uber, Fourth Quarter and Full Year 2025 Results and prepared remarks, February 2026: 3,751 million trips in the quarter; "40 million+ trips per day"; 9.7 million monthly drivers and couriers; 202 million monthly active platform consumers.
- Ranney, Scaling Uber's Real-time Market Platform, QCon London, March 2015: the state digest on drivers' phones.
- AWS documentation: Aurora DSQL, its resilience and quotas (checked September 2026).
Uber hasn't published its fleet sizes, its failover times, whether Docstore's cross-region replication is synchronous, or how state digests are versioned; those parts of this round, the site ownership controls and the whole AWS translation are ours.
R3.6 Numbers and Cost
Figures marked published come from the sources above; everything else is our assumption or derived from one.
Scale today (published, Q4 2025)
| Quantity | Arithmetic | Result |
|---|---|---|
| Trips in the quarter | Published | 3,751 million |
| Trips a day | 3,751M ÷ 92 days | ≈ 40.8M ("40 million+", published) |
| Average rate | 40.8M ÷ 86,400 | ≈ 472/s |
| Peak rate | × 2.5 (assumption: a global business flattens the peak) | ≈ 1,180/s |
| Monthly drivers and couriers | Published | 9.7M |
| Online at the global peak | Assumption | 3M |
| Pings | 3M ÷ 4 s | 750,000/s |
Why the fulfillment rewrite paid for itself (efficiency)
| Arithmetic | Cores for 3M online | |
|---|---|---|
| Old stack: 20 online drivers/couriers per core (published, 2021) | 3,000,000 ÷ 20 | 150,000 |
| A new stack at 200 per core (our assumption; Uber hasn't published its new ratio) | 3,000,000 ÷ 200 | 15,000 |
Capacity per site (our design). Each of two regions must carry 100% of the load if the other fails: 15,000 cores per region, 5,000 per AZ in three AZs, 30,000 in all, half of it idle-ready at any moment. On Arm instances one vCPU is one core, so each region needs a Running On-Demand Standard instances quota of at least 15,000 vCPUs, plus about 10% for the extra instances a rolling deploy adds: 16,500 vCPUs per Region, checked in both Regions before we need them.
H3 (from step 3.1)
| Arithmetic | Result | |
|---|---|---|
| Resolution 8 cell | Published average | 0.737 km², edge 0.531 km |
| Distance between neighboring centers | √3 × 0.531 | ≈ 0.920 km |
| Cells for a 2 km search | k = 3: 1 + 3 × 3 × 4 | 37 cells, ≈ 27.3 km² |
| Cells in a 600 km² city | 600 ÷ 0.737; at resolution 7: 600 ÷ 5.16 | ≈ 814; ≈ 116 |
Docstore commits (from step 3.2; assumed 1 ms between zones). A 3-node partition commits when the leader and one follower have the transaction: about one cross-zone round trip plus a log flush, a millisecond or two. A 5-node partition waits for the second-fastest of four followers.
Region failover (from step 3.4)
| Chain | Minutes | Share of a 4.32-minute monthly budget | |
|---|---|---|---|
| Human-gated | 30 + 180 + 5 + 60 + 30 + 1 s | 5.1 | 118% |
| Pre-approved for clear signals | 30 + 10 + 5 + 60 + 30 + 1 s | 2.3 | 52% |
Transfer costs in the AWS translation (list prices; $0.02 per GB across AZs, $0.02 per GB between us-east-1 and us-west-2; 30-day month)
| Stream | Arithmetic | A month |
|---|---|---|
| Pings: forward to owner and copy to 2 other AZs | 750,000 × 300 B × (2/3 + 2) = 600 MB/s = 51,840 GB a day × 30 × $0.02 | $31,104 |
| Pings copied to the other Region (we don't) | 750,000 × 300 B = 225 MB/s = 19,440 GB a day × 30 × $0.02 | $11,664 avoided |
| Trip record copied to the other Region | 40.8M trips × 30 KB (assumption) = 1,224 GB a day × 30 × $0.02 | ≈ $734 |
The expensive bytes are the fleeting ones. The permanent record is cheap to copy anywhere; the ping stream is what a design must keep inside a zone where it can.
Published cost shape of the Spanner move: nodes about 80% of the cost, network transfer about 20% (because the applications ran on-premises), and multi-region configurations 3–4× the price of single-region ones. Storage was "trivial" because only live data stays.
R3.7 Trade-Offs
Hexagons vs squares
| Squares (geohash, S2 faces) | Hexagons (H3) | |
|---|---|---|
| Neighbor distances | Two (edge and corner) | One |
| Exact nesting | Yes: four children fill the parent | No: seven children only approximately fill it |
| Radius search | A square ring of cells | A k-ring, closer to a circle |
| Pentagons or seams | Face edges (S2), distortion (geohash on Mercator) | 12 pentagons, in the ocean |
| Best for | Exact hierarchies, range scans on a curve | Movement, smoothing, comparing neighbors |
Build vs buy datastores, over the whole loop
| System | Built or bought | Why, in Uber's words or ours |
|---|---|---|
| Schemaless (2014) | Built | "Operational trust"; nothing met all five requirements in time |
| Cassandra (later) | Adopted | Flexibility; later found "lacking operational maturity at Uber's scale" |
| Docstore (2021) | Built, from Schemaless | A transactional store on the engine the team already ran |
| Spanner for fulfillment (2021) | Bought | Transactions across regions with "low operational overhead" |
The pattern: build where the team's operational knowledge is the advantage and the workload is Uber-specific; buy where the hard part (global consensus, TrueTime) is someone else's core business.
An AP ring vs a CP transactional store for live state
| Ringpop and sagas (2014) | Transactions (2021) | |
|---|---|---|
| During a partition | Keeps serving; may diverge | Minority side stops writing |
| Multi-entity change | Saga with compensations | One transaction |
| Debugging | Hard: "difficult to debug issues in production" (Uber) | State is what was committed |
| Latency | In-memory, one hop | A commit across regions |
| Scale limit | Ring size per pod | The store's sharding |
Active-active vs active-passive. Active-active keeps both regions warm and proven every day, at the price of standby capacity inside every region and an ownership protocol per city. Active-passive is cheaper to run and more likely to fail on the day it's used.
What changed from Round 1. Round 1 split one process and one database so each could grow. Round 2 made growth automatic by giving up consistency: gossip, last-write-wins, sagas and eventually consistent indexes. Round 3 bought consistency back where it matters (who has which trip, one transaction per accept) and kept the append-only, log-driven ideas where they fit (the record store, CDC, post-commit tasks). The lesson of the loop: at each scale, choose the weakest guarantee the business can live with, and be ready to buy stronger ones back when debugging the weak ones costs more than the stronger store.
R3.8 Failure Modes
| Trigger | What you'd see | How the design responds |
|---|---|---|
| A region fails mid-trip | Its cities' apps lose their connection | Outside detection, gated move, epoch bump, phones resync; live state is already in the survivor (step 3.4). |
| The old region comes back | Stale services try to write for cities they no longer own | Every write carries owner_epoch; the store and the services refuse lower epochs. |
| A network partition inside a Raft partition | The minority's leader can't commit; writes for that partition stall briefly | The majority side elects a leader and continues; clients retry against the new leader. |
| A zone fails | One node of every 3-node partition is gone | Each partition still has 2 of 3. Our 5,000 cores per AZ assume the other region absorbs the lost zone's cities; to ride out a zone loss inside one region alone, each AZ needs 50% more (7,500). |
| Transaction contention | ABORTED transactions and retries on hot rows (a busy merchant, a driver with many offers) | Keep transactions small; retry with jitter; route a hot entity's work through one owner so it doesn't contend with itself. |
| A LATE task runs twice | A duplicate notification | At-least-once by design; the task carries an idempotency key. |
| A snapshot misses the last change | A fare adjustment disappears after the trip is deleted from the live store | Snapshot from a strong read; delete only if the version hasn't changed. |
| A cloud interconnect fails | Higher latency, then errors, on calls to the cloud database | N+1 links at two layers (published); alarm on the redundancy, not just on the traffic. |
| A stream consumer hits a poison change | Flux or LATE consumers stop advancing on one shard | Retry, then move the record aside; bisect a failing batch; refill from the source of truth (Docstore itself), not by replaying the raw stream. |
R3.9 Runbook and Incident Response
| Signal | Alarm | Severity | First action |
|---|---|---|---|
| Time from request to first offer, p99 per city | > 2× its usual value for 5 min | P1 | Which cities? One pod, one region, or everywhere? |
| Ring membership churn (where Ringpop remains) | Faulty transitions > 5 a minute, or more than one membership checksum for 1 min | P2 | Common cause first: a deploy, a network device, GC |
| Shard write latency, p99 per partition | > 3× baseline for 5 min | P2 | Hot partition or a leader change? |
| Buffer or retry queue depth | Growing for 10 min | P2 | A primary is down or slow |
| Trigger and CDC lag per shard | > 60 s | P2 | A poison record? Too few consumers? |
| Transaction abort rate | > 2% for 10 min | P3 | Which rows contend? |
| Outside probes per region | 3 consecutive failures | P1 | Start the region runbook |
| Standby headroom | Survivor can't carry 100% of the peak | P2 | Raise capacity or quota before it's needed |
Procedure: ring instability REL 11
- Stop changing membership. Pause deploys and autoscaling for the ring's service.
- Find the common cause. Are the flapping nodes in one zone, one rack, one deploy batch? Are they in long GC pauses?
- Remove the sick, not the noisy. Take nodes with confirmed problems out deliberately, a few at a time, letting the ring converge (checksums agree) between each.
- Only then tune. If healthy nodes are being suspected, lengthen the suspicion window before shortening anything.
Procedure: shard or partition overload PERF 3
- Confirm it's one partition: latency is high for one partition's keys, normal elsewhere.
- Find the key: a hot trip, merchant or city? A key can't be split by adding nodes.
- Shed or reshape: rate-limit the source; move other shards off that partition; for a structurally hot key, change its partition key.
Procedure: region evacuation REL 13 · OPS 10
- Confirm from outside that the region, not a service, is down.
- Check the survivor has capacity and quota for 100% of the peak, and may hold the affected cities' data.
- Flip routing controls for the affected cities; confirm each city's epoch moved.
- Watch resyncs and the first-offer latency in the moved cities.
- Fence the failed region's services before it rejoins; move cities back one at a time.
Go deeper: CLI and SQL checks (AWS translation; replace names with real ones)
text# 1. Which targets behind the driver connection tier are unhealthy aws elbv2 describe-target-health --target-group-arn arn:aws:elasticloadbalancing:us-east-1:111122223333:targetgroup/driver-conn/abc123 # 2. Scale the dispatch group in one AZ aws autoscaling set-desired-capacity --auto-scaling-group-name dispatch-use1-az1 --desired-capacity 120 # 3. Read, then flip, a region's routing control (run against a cluster endpoint) aws route53-recovery-cluster get-routing-control-state --routing-control-arn arn:aws:route53-recovery-control::111122223333:controlpanel/abc/routingcontrol/def --region us-west-2 --endpoint-url https://host-aaa.us-west-2.example.com/v1 aws route53-recovery-cluster update-routing-control-state --routing-control-arn arn:aws:route53-recovery-control::111122223333:controlpanel/abc/routingcontrol/def --routing-control-state Off --region us-west-2 --endpoint-url https://host-aaa.us-west-2.example.com/v1 # 4. The vCPU quota for Running On-Demand Standard instances, in the survivor Region aws service-quotas get-service-quota --service-code ec2 --quota-code L-1216C47A --region us-west-2 # 5. Alarms currently firing for dispatch aws cloudwatch describe-alarms --state-value ALARM --alarm-name-prefix dispatch- # 6. On a MySQL replica: how far behind is it? SHOW REPLICA STATUS\G # 7. On a MySQL master of the record store: how many buffered cells wait for replay? SELECT COUNT(*) FROM buffered_cells;
R3.10 Pillar Check
| Pillar | What Round 3 adds |
|---|---|
| Reliability | Synchronous live state across regions; outside detection and a gated, epoch-fenced city move; Raft partitions across zones; standby capacity for a full peak REL 13 · REL 10 · REL 1 |
| Performance Efficiency | One hexagonal grid for every team; transactions coalesced to cut round trips; the old stack's 20 drivers per core replaced PERF 1 · PERF 3 |
| Security | Location history classified as sensitive and kept in the live store only while the trip is live; state digests on phones are encrypted (published) SEC 7 · SEC 8; traffic between Uber's sites and the cloud database travels over dedicated interconnects, with TLS on every connection (our design) SEC 9 |
| Cost Optimization | The cost split of a multi-region database understood (80% nodes, 20% network); pings kept out of cross-region replication; multi-region only for the state that needs it COST 8 · COST 5 · COST 11 |
| Operational Excellence | Per-city signals; runbooks for ring churn, hot partitions and region evacuation; transaction analysis for aborts OPS 8 · OPS 10 · OPS 11 |
| Sustainability | A tenfold efficiency target over the old 20 drivers per core (our target, not a published figure); live data deleted once the trip is snapshotted; MyRocks cutting the disk footprint SUS 3 · SUS 4 |
R3.11 Round 3 Rubric and Follow-Ups
What an architect (L7) answer adds over L6
- Chooses consistency per kind of data: synchronous for contended live state, asynchronous for the record and analytics, none for pings.
- Explains why a hexagonal grid suits a marketplace, and computes a k-ring.
- Compares an AP ring with sagas against a transactional store, using Uber's own published reasons.
- Sizes a Raft partition (3 vs 5) and says what a minority can do.
- Designs a region failover end to end: outside detection, a gated move, capacity and residency in the survivor, epochs, resync, and a timing chain checked against the error budget.
- Separates what Uber published from what it designs, and reads published numbers for what they imply.
Follow-up questions
-
"Why did Uber buy Spanner for fulfillment but build Docstore for everything else?" Answer: fulfillment needed multi-row, multi-table transactions across regions, the hardest thing a database does, and wanted "low operational overhead"; Spanner provides exactly that. Docstore serves many teams on Uber's own infrastructure, where Uber already had deep MySQL expertise and wanted control of cost and features. Different workloads, different make-or-buy answers.
-
"What's the RPO of your design if region A fails?" Answer: zero for live fulfillment state, which commits synchronously in both regions. For asynchronously replicated data, its replication delay at the moment of failure, typically seconds. Pings aren't replicated: they refill in 4 seconds.
-
"Could you make the failover fully automatic?" Answer: for unambiguous signals, yes: probes from two outside places and the surviving region all agree. The risk is a partition that makes a healthy region look dead; the epoch fence makes that safe for data, but not free for riders, who'd see a needless move. So automate the clear cases and page a human for the rest.
-
"Is Ringpop still how Uber dispatches?" Answer: Uber's 2021 posts describe the fulfillment platform moving off its Ringpop-based stack to Spanner, and the Node.js library's repository says it "is no longer under active development". The Go library, ringpop-go, is still on GitHub and is a dependency of Uber's open-source Cadence workflow engine. The ideas (a hash ring, SWIM, handle-or-forward) remain standard.
Interview gotchas from this round's wrong answers
| Gotcha | Why it's wrong |
|---|---|
| "A mapping table between everyone's zones" | Every pair needs a mapping, and every change breaks history. |
| "2PC on top of an append-only store" | All of 2PC's blocking failures, on a design built to avoid coordination. |
| "Global tables will keep accepts consistent" | Local-replica conditions and last writer wins let two regions both accept. |
| "Replicate everything to the other region" | Pays most for the data that matters least. |
| "Health checks inside the failed region decide" | They fail with it. |
| "Failover is automatic, so it's instant" | Detection, approval, DNS and backoff add up to minutes. |
Loop Closer: Interview Strategy for All Three Rounds
How to Run Each 60-Minute Round
| Time | Round 1 | Round 2 | Round 3 |
|---|---|---|---|
| 0–5 min | Scope: request, match, trip, pay; what's breaking | Restate Round 1 in 60 seconds | Restate Round 2 in 60 seconds |
| 5–15 min | Requirements; pings and trip API; trip states | Scope raise → what breaks | Scope raise → what breaks |
| 15–40 min | Steps 1.0–1.3: dispatch by city and a grid → trip store out of the shared database → services | Steps 2.1–2.7: ring → SWIM → cells → buffered writes → triggers → indexes → fencing | Steps 3.1–3.5: H3 → Docstore → transactions for live state → site loss → cloud |
| 40–50 min | Dispatch load, the runway under compound growth | Trips and cells, dispatch fleet, pause tails, NLB and cross-AZ costs, trigger polling | Scale, efficiency, capacity per site, failover chain vs budget, transfer costs |
| 50–60 min | Failures and pillar check | Failures, gotchas, pillar check | Failures, runbook, pillar check |
For how to spend a single 45-minute round, see the 45-minute interview blueprint. For offers, accepts, batching and surge in a generic dispatch service, see the ride-sharing dispatch loop; for grids and radius search, the proximity loop; for another store that outgrew its first design, the Discord case study.
The Two Sentences That Matter Most
- Opening any round: "Pings are fleeting and huge, trips are permanent and valuable, so I'll keep positions in memory, sharded by map cell, and put trips in a store that grows by adding shards and never loses an acknowledged write."
- When scale arrives: "I'll pick the weakest guarantee each kind of data can live with: no replication for pings, an append-only log with at-least-once consumers for the trip record, and synchronous transactions only for who has which trip, with an epoch on every write so a recovering site can't overwrite the new owner."
Well-Architected Review Sheet
Interviewers rarely ask "which pillar is this?". They ask the pillar's question in plain words. Rehearse one sentence per row.
| Pillar | Question you'll hear | One-sentence answer | Round | Backed by |
|---|---|---|---|---|
| Reliability | "What happens when a dispatch machine dies?" (REL 11) | SWIM marks it faulty in about 13 seconds, the ring moves its cells, and pings refill them within 4 seconds. | 2 | Step 2.2 |
| "How do you stop one city from taking down the others?" (REL 10) | Shard by city, then by map cell, and group cities into pods so a failure has a small blast radius. | 1–3 | Steps 1.1, 2.1, R3.1 | |
| "How do you change the database under live traffic?" (REL 8) | One library in front of the store, mirrored writes, a backfill that writes the same versions, and sampled validation before the switch. | 1–2 | Steps 1.2, 2.3 | |
| "What if a whole region fails?" (REL 13) | Live state is already synchronous in both regions; outside probes, a gated move, an epoch bump and phone resync bring cities back in minutes. | 3 | Step 3.4 | |
| Performance | "How do you find nearby drivers fast?" (PERF 3) | Keep positions in memory, keyed by map cell, and search only the cells around the pickup: a 3-ring of 37 hexagons for 2 km. | 1–3 | Steps 1.1, 2.1, 3.1 |
| "Why is your p99 worse than your p50?" (PERF 2) | A search touches eight processes, so 3.9% meet a pause; backup requests to a second replica cut that tail. | 2 | R2.6 | |
| Security | "Location data is sensitive. How do you treat it?" (SEC 7) | Pings live only in memory; trips are classified sensitive and live data leaves the live store once snapshotted. | 1–3 | R1.10, R2.10, R3.10 |
| Cost | "Where does transfer cost hide?" (COST 8) | In the ping stream across zones (about $31,000 a month at 750,000 pings a second in our translation), not in copying trips between regions (about $734). | 2–3 | R2.6, R3.6 |
| "Build or buy the database?" (COST 11) | Build where the team's operating knowledge is the edge (Schemaless, Docstore); buy where global consensus is the hard part (Spanner). | 2–3 | R2.7, R3.7 | |
| Operations | "How do you know dispatch is healthy?" (OPS 8) | Time to first offer per city, ring churn, partition latency, buffer and trigger lag, abort rate, outside probes. | 2–3 | R2.8, R3.9 |
| Sustainability | "How did you do more with less?" (SUS 3) | The Go rewrite cut Schemaless CPU by more than 85%, and the fulfillment rewrite replaced a stack that managed only 20 online drivers per core. | 2–3 | R2.10, R3.6 |
Rubric Across Levels
| Dimension | L5 (Round 1) | L6 (Round 2) | L7 (Round 3) |
|---|---|---|---|
| Dispatch | Shards by city; a grid instead of a scan | Consistent hashing with replica points; SWIM; handle-or-forward | One hexagonal grid; transactions for live state; pods and regions |
| Trip storage | Gets trips out of the shared database in time | Append-only cells, buffered writes, idempotent retries | Transactions per partition with Raft; snapshot-then-delete done safely |
| Change propagation | Knows the synchronous chain is fragile | Triggers from the store's own log; poison cells; polling vs log tailing | Outbox in the same transaction; why commit-timestamp scans are safe on Spanner and not on MySQL |
| Failure | A process crash; a full disk | Flap storms, master loss, split brain fenced at the store | Region loss with outside detection, epochs, capacity, residency and a checked timing chain |
| Honesty | Labels assumptions | Separates Uber's published mechanisms from its own fixes | Reads Uber's numbers for what they imply and says what they don't show |
Sources
All sources used on this page, oldest first.
- Ranney, Scaling Uber's Real-time Market Platform, QCon London, March 2015, and the High Scalability summary of that talk.
- Schmidt, Project Mezzanine: The Great Migration, Uber blog, July 2015.
- Haddad, Service-Oriented Architecture: Scaling the Uber Engineering Codebase As We Grow, Uber blog, September 2015.
- Thomsen, Designing Schemaless, The Architecture of Schemaless and Using Triggers On Schemaless, Uber blog, January 2016.
- Lozinski, How Ringpop from Uber Engineering Helps Distribute Your Application, Uber blog, February 2016.
- Reuters, Uber reaches 2 billion rides six months after hitting its first billion, July 2016.
- Klitzke, Why Uber Engineering Switched from Postgres to MySQL, Uber blog, July 2016.
- Nielsen and Johnsen, Rewriting the Sharding Layer of Uber's Schemaless Datastore, Uber blog, February 2018.
- Brodsky, H3: Uber's Hexagonal Hierarchical Spatial Index, Uber blog, June 2018.
- Chatterjee, Chaudhary and Tariq, Evolving Schemaless into a Distributed SQL Database, Uber blog, February 2021.
- Neerabail, Medisetty, Thangavelu and others, Uber's Fulfillment Platform: Ground-up Re-architecture, Uber blog, July 2021.
- He, Medisetty and others, Building Uber's Fulfillment Platform for Planet-Scale using Google Cloud Spanner, Uber blog, September 2021.
- Uber, MySQL to MyRocks Migration in Uber's Distributed Datastores, Uber blog, September 2022.
- Pozniansky and others, How Uber Serves Over 40 Million Reads Per Second from Online Storage Using an Integrated Cache, Uber blog, February 2024.
- Uber, Adopting Arm at Scale: Bootstrapping Infrastructure, Uber blog, February 2025.
- Uber, Fourth Quarter and Full Year 2025 Results, February 2026.
- Uber open source: ringpop-node, ringpop-go, H3 and Cadence (checked September 2026).
- H3, Tables of Cell Statistics Across Resolutions; S2 Geometry, S2 Cell Statistics (checked September 2026).
- MySQL 8.4 Reference Manual, Semisynchronous Replication.
- AWS, Aurora DSQL documentation, Elastic Load Balancing pricing and EC2 data transfer pricing (
us-east-1, checked September 2026), for the translation's figures.