Design Uber's Real-Time Dispatch & Schemaless MySQL Architecture
1. Problem Statement & Scope Clarification
System Mission
Design Uber’s real-time geospatial dispatch platform and distributed append-only storage engine (Schemaless). The system must match millions of riders with nearby drivers in real time ( matching latency), ingest continuous GPS streams from over active drivers ( location updates/sec), partition stateful dispatch services via a decentralized gossip-based consistent hashing ring (Ringpop), and persist billions of immutable trip lifecycle transitions across a horizontally sharded MySQL InnoDB datastore with zero in-place updates.
The Schemaless Architectural Invariant: Schemaless is an append-only datastore layered on top of horizontally sharded MySQL InnoDB clusters. Data is organized into cells keyed by (Row Key, Column Name, Ref Key). Cells are strictly immutable. Updates never overwrite existing rows; instead, they append a new cell with an incremented ref_key (version). An auto-incrementing added_id primary key provides a strictly ordered, monotonically increasing changelog that enables reliable binlog tailing and Change Data Capture (CDC) without distributed locking.
Synthesizing vector architecture diagram...
Functional Requirements
- Real-Time Driver Location Ingestion: Ingest high-frequency GPS telemetry () from drivers every 4 seconds, indexing locations into geospatial hierarchical cells (Google S2 / Uber H3).
- Geospatial Dispatch Matching (DISCO): When a rider requests a trip, identify the top optimal available drivers within a dynamic geospatial radius based on Estimated Time of Arrival (ETA), vehicle category, and supply-demand balance in .
- Decentralized Request Routing (Ringpop): Distribute active trips and supply state across stateful dispatch worker nodes using consistent hashing with a SWIM gossip protocol, automatically re-balancing ring membership upon node failure.
- Append-Only Immutable State Persistence (Schemaless): Record every state transition of a trip (
REQUESTED,OFFERED,ACCEPTED,ARRIVING,IN_TRANSIT,COMPLETED,BILLED) as an immutable cell in Schemaless. - Asynchronous Secondary Indexing: Decouple query indexing from write paths. Tail the MySQL binlogs via Kafka to maintain secondary indexes (e.g., query trips by
rider_uuid,driver_uuid, orcreated_at_time) in Amazon OpenSearch and relational index tables.
Non-Functional Requirements (SLAs & SLOs)
- High Availability: availability for driver-rider matching and ride state updates ( minutes downtime per year).
- Matching Latency: , , for dispatch match completion.
- GPS Ingestion Throughput: Ingest with write latency to in-memory spatial caches.
- Storage Durability: Zero data loss () for confirmed trip state cells across multi-AZ MySQL InnoDB replicas with semi-synchronous replication.
- CAP / PACELC Classification:
- Real-Time Dispatch (Ringpop / DISCO): AP system under CAP; PA/EL under PACELC. Availability and low latency over strict consistency; localized temporary supply mismatches are resolved via optimistic locking rather than blocking distributed locks.
- Schemaless Storage & Billing: CP system enforcing strict ACID guarantees per shard partition key (
row_key = trip_uuid).
Out-of-Scope
- Turn-by-turn routing navigation algorithms and street-level map graph rendering (handled by Routing Engine).
- Fare calculation machine learning pricing models (Surge pricing rates ingested as precomputed inputs).
- Regulatory tax reporting and local municipal compliance filings.
2. Capacity & Scale Estimation (Back-of-the-Envelope Math)
Traffic Scale Calculations
- Active Riders: Monthly Active Users (MAU); Daily Active Users (DAU).
- Active Drivers: drivers globally; concurrently online during peak hours.
- Completed Trips per Day: trips/day.
- Driver GPS Ingestion QPS: Each online driver pings location every 4 seconds (): (All 6 M registered drivers are never online at once; the 1.5 M concurrent figure above is what the fleet, Redis TTLs and DISCO sizing are built on.)
- Dispatch Ride Request QPS:
- State Transition Writes to Schemaless: Each trip generates an average of immutable state transitions (events, location pings, payment auths, status updates):
Storage Calculations (3-Year Horizon)
- Schemaless Cell Size Breakdown:
added_id: 8 bytes (BIGINT auto-increment)row_key: 16 bytes (UUID binary)column_name: 32 bytes (VARCHAR)ref_key: 8 bytes (BIGINT version)body(JSON/MsgPack compressed payload): 1,200 bytes- Indexes and InnoDB row overhead: 150 bytes
- Total per cell: .
- Daily Raw Storage Growth:
- 3-Year Raw Storage Volume:
- Replication Multiplier ( Across Multi-AZ MySQL Shards):
Memory & Cache Sizing (Driver Locations & DISCO Ring)
- Active Driver Spatial Index (H3 Resolution 8 Hexagons): (fits entirely in memory).
- Active Trip Contexts: concurrent in-progress trips .
- ElastiCache Redis Spatial Tier (Driver Geofence Index): Provisioned across 6 nodes in cluster mode: for working sets and sub-millisecond proximity queries.
Fleet Sizing & Compute Resources
- DISCO Dispatch Workers (Ringpop Nodes on AWS ECS Fargate - 4 vCPU, 16 GB RAM): Each node coordinates active drivers and processes dispatch match evaluations/sec.
- Schemaless MySQL Database Shards: Dividing across database shards (each AWS Aurora / MySQL node comfortably handles ):
3. AWS-First High-Level Architecture
Synthesizing vector architecture diagram...
Follow a ride request from the top. Rider and driver apps enter through Route 53 and the NLB to any DISCO dispatch node. In the "Real-Time Dispatch Tier" panel, the nodes form a Ringpop ring kept in sync by gossip: each H3 area is owned by exactly one node, so any node that receives a request forwards it to the owner, which looks up drivers in the Redis supply grid and makes the match. In the "Schemaless Storage Tier" panel, trips are written through the Schemaless router to about 20 MySQL shards, each with a read replica. In the "Change Data Capture & Async Processing Tier" panel, each shard's binlog streams into Kafka, where consumers build secondary indexes in OpenSearch and archive to S3. Owning each area on exactly one node means matching needs no distributed locks between nodes.
Data Flow Walkthrough
- Driver GPS Ingestion Flow:
- Drivers emit GPS pings every 4 seconds. The request hits AWS NLB and routes to a DISCO worker node.
- The DISCO worker maps the driver's coordinates to an Uber H3 hexagonal cell index (resolution 8, edge length , about across).
- The driver's location, heading, and status (
AVAILABLE) update in ElastiCache Redis with a 10-second TTL.
- Rider Dispatch Request Flow:
- A rider requests a pickup at coordinate . The request reaches a DISCO node.
- The receiving node consults the Ringpop consistent hash ring. Ringpop hashes the trip's geographic H3 cell or rider UUID to determine the authoritative coordinator node responsible for this geographic partition. If the current node is not the owner, it proxies the request to the target DISCO node in .
- The coordinator DISCO node queries Redis for all
AVAILABLEdrivers located within the pickup H3 hexagon and neighboring concentric k-rings. - Drivers are scored via a batch ETA matrix, and the top candidate is selected.
- Schemaless Append-Only Write Flow:
- The match generates a trip record:
row_key = trip_uuid. - The DISCO service writes an immutable cell to the Schemaless Proxy:
column_name = "status",ref_key = 1,body = {"state": "REQUESTED", "rider_id": "...", "origin": {...}}. - The Schemaless Proxy hashes
row_keyacross 20 MySQL shards:shard_id = MurmurHash3(row_key) % 20. - The target MySQL primary executes a single fast append into InnoDB:
INSERT INTO schemaless_cells (row_key, column_name, ref_key, body) VALUES (...). Zero lock contention occurs because the write is an append, not an in-place update.
- The match generates a trip record:
- CDC Secondary Index Pipeline:
- A binlog tailing agent (Debezium/Maxwell) captures the newly appended cell from MySQL's binary log and publishes the change event to Amazon MSK (Kafka).
- Async index workers consume the Kafka topic, extract secondary attributes (
rider_id,driver_id,created_at), and update the secondary query index in Amazon OpenSearch, allowing support agents and users to query historical trips.
End-to-End Request Tracing Walkthrough
| Step # | Event / Action | Component State | Distributed Transition | Output / Response |
|---|---|---|---|---|
| Step 1 | Driver App emits GPS location ping | Driver mobile socket | TCP connection terminated at AWS NLB; payload forwarded to DISCO node | Coordinates (37.7749, -122.4194) received |
| Step 2 | Geospatial H3 Cell Encoding | DISCO Node processing | Coordinates converted to 64-bit H3 index: 8828308281fffff | Spatial cell resolution 8 assigned |
| Step 3 | Driver Ephemeral State Commit | DISCO ElastiCache | Atomic Redis pipeline: SADD cell:8828308281fffff d109, SETEX driver:d109 10 AVAILABLE | Driver marked available in spatial grid () |
| Step 4 | Rider invokes RequestDispatch API | Rider mobile client | Rider request arrives at Ingress NLB; routed to DISCO Node 42 | HTTP POST /v1/dispatch/trips received |
| Step 5 | Ringpop Consistent Hash Routing | DISCO Node 42 active | Ringpop evaluates MurmurHash(pickup_h3); determines Node 15 is coordinator | Node 42 proxies request internally to Node 15 () |
| Step 6 | Proximity Radius Driver Discovery | DISCO Node 15 coordinating | Node 15 queries Redis for available drivers in primary cell and 1-ring neighbors (7 cells, guaranteed reach ) | Returns 12 candidates within ~0.75 miles |
| Step 7 | ETA Routing & Candidate Scoring | DISCO Match Engine | Candidates ranked based on routing network distance and traffic vectors | Driver d109 selected with lowest ETA () |
| Step 8 | Schemaless Cell Append (REQUESTED) | DISCO Schemaless Proxy | Router hashes trip_uuid; routes to MySQL Shard 07 Primary | INSERT INTO schemaless_cells commits in |
| Step 9 | Driver Push Notification Dispatched | Dispatch Engine | Apple APNs / FCM push dispatched to Driver d109 with 15-second acceptance lease | Driver mobile screen rings with offer modal |
| Step 10 | Driver Accepts Trip Offer | Driver App taps Accept | Driver submits acceptance token to DISCO Node 15 | DISCO verifies lease validity () |
| Step 11 | Schemaless Cell Append (ACCEPTED) | DISCO Schemaless Proxy | Appends new cell: column = "status", ref_key = 2, state = ACCEPTED | Immutable cell appended; previous cell preserved |
| Step 12 | Binlog CDC Secondary Indexing | MySQL Shard 07 Binlog | Debezium captures append; emits event to Kafka schemaless.trips topic | OpenSearch secondary index updated in |
4. API Interface Design & Wire Protocol
1. Trip Dispatch Request (POST /v1/dispatch/trips)
httpPOST /v1/dispatch/trips HTTP/1.1 Host: dispatch.uber.internal Content-Type: application/json X-Client-UUID: rider_01HZX894KMNPQ Idempotency-Key: 9a8b7c6d-5e4f-3a2b-1c0d-9e8f7a6b5c4d { "rider_uuid": "rider_01HZX894KMNPQ", "pickup_location": { "latitude": 37.774929, "longitude": -122.419418, "h3_index": "8828308281fffff", "street_address": "1455 Market St, San Francisco, CA" }, "dropoff_location": { "latitude": 37.789172, "longitude": -122.401449, "h3_index": "8828308295fffff", "street_address": "Montgomery St Station, San Francisco, CA" }, "vehicle_category": "UBER_X", "surge_multiplier": 1.2, "max_acceptable_eta_seconds": 600 }
Response: 201 Created
json{ "trip_uuid": "trip_84920194-7f12-4000-8000-abcdef123456", "status": "DISPATCHING", "estimated_pickup_epoch_ms": 1773648192000, "surge_multiplier": 1.2, "assigned_driver": null, "created_at_epoch_ms": 1773648000000, "poll_interval_ms": 2000 }
2. Driver Location Ingestion Stream Protocol (gRPC / Protobuf)
protobufsyntax = "proto3"; package uber.dispatch.v1; service LocationIngestionService { rpc PushDriverLocation (DriverLocationRequest) returns (DriverLocationResponse); rpc BatchDriverLocations (BatchDriverLocationRequest) returns (DriverLocationResponse); } message DriverLocationRequest { string driver_uuid = 1; double latitude = 2; double longitude = 3; float bearing_degrees = 4; float speed_meters_per_sec = 5; uint64 h3_index = 6; string vehicle_category = 7; enum DriverStatus { AVAILABLE = 0; ASSIGNED = 1; IN_TRANSIT = 2; OFFLINE = 3; } DriverStatus status = 8; int64 timestamp_epoch_ms = 9; } message DriverLocationResponse { bool acknowledged = 1; int32 next_reporting_interval_seconds = 2; }
3. Status Codes & Error Contracts
| HTTP Status | Reason Code | Error Contract Payload | Client Handling / Recovery Strategy |
|---|---|---|---|
201 Created | DISPATCH_INITIATED | Standard trip response with trip_uuid | Client displays "Finding your ride" radar animation |
400 Bad Request | INVALID_GEOPOLYGON | {"error": "Coordinates outside operating territory"} | Alert user that service is unavailable in this territory |
409 Conflict | TRIP_ALREADY_ACTIVE | {"error": "Active trip exists", "trip_uuid": "..."} | Redirect client to live active trip screen |
422 Unprocessable | NO_DRIVERS_AVAILABLE | {"error": "Zero supply in geospatial cell"} | Expand search radius; prompt user to retry in 2 minutes |
429 Too Many Requests | DISPATCH_RATE_LIMIT | {"error": "Too many requests", "retry_after": 1.0} | Debounce user requests on client device |
503 Service Unavailable | RINGPOP_RESHARDING | {"error": "Gossip partition rebalancing"} | Retry request; Ringpop automatically forwards to new owner |
5. Data Models & Storage Architecture
The Schemaless MySQL Relational Schema
Schemaless structures data into shards, rows, and cells. A single row represents a business entity (e.g., a trip). Within a row, data is organized into columns (e.g., status, driver, payment, rating). Each column contains multiple cells representing the historical timeline of updates, keyed by ref_key.
Synthesizing vector architecture diagram...
Read the three tables as layers. A SCHEMALESS_ROW is one entity (such as a trip), identified by row_key, which also decides its shard. Its data lives in SCHEMALESS_CELLs: each cell holds one column (such as "status" or "payment") at one version (ref_key) as a JSON blob, and cells are only ever appended, never updated in place, so the full history of every column is kept and concurrent writers never overwrite each other. added_id gives every cell a global order, which is what change streams read. SECONDARY_INDEX maps another key (such as driver ID) back to a row_key, and is built asynchronously from those streams. This append-only cell design is what let Uber use sharded MySQL as a scalable, schemaless store.
Production MySQL DDL for Schemaless Cells Table
sql-- Executed on each of the 20 MySQL Shards CREATE TABLE schemaless_cells ( added_id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT, row_key VARBINARY(64) NOT NULL, column_name VARCHAR(64) NOT NULL, ref_key BIGINT UNSIGNED NOT NULL, body MEDIUMBLOB NOT NULL, -- Compressed JSON / MsgPack payload is_deleted TINYINT UNSIGNED NOT NULL DEFAULT 0, created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, PRIMARY KEY (added_id), -- Clustered Index: Enables sequential binlog tailing UNIQUE KEY uq_cell (row_key, column_name, ref_key) -- Deduplication & Idempotency; its (row_key, column_name) prefix also serves -- "latest state" lookups (ORDER BY ref_key DESC LIMIT 1), so no separate index is needed ) ENGINE=InnoDB DEFAULT CHARSET=binary ROW_FORMAT=COMPRESSED KEY_BLOCK_SIZE=8;
Why added_idis the Clustered Primary Key
In traditional relational databases, the primary key would be (row_key, column_name, ref_key). However, in a distributed append-only system, using added_id AUTO_INCREMENT as the clustered primary key provides two critical production advantages:
- Zero B-Tree Page Splits on Writes: Every new write appends strictly to the end of the InnoDB clustered index leaf pages. Disk writes are sequential, eliminating costly page splits and minimizing disk I/O.
- Deterministic, Monotonic Binlog Tailing: CDC tailing engines (Debezium/Maxwell) checkpoint their position simply by storing the highest
added_idobserved. Catching up after worker restarts requires a simple query:SELECT * FROM schemaless_cells WHERE added_id > :last_checkpoint ORDER BY added_id ASC LIMIT 1000.
Unlock Complete Architecture & Production Runbooks
You have explored the free architectural preview (~45%). Spend 1 Coin to unlock the remaining 6 production deep-dive sections for a full 24 hours.