Skip to content
Reliable Data Engineering
Practice problem hard geospatialstreamingh3flinkfeature-store
Practise with timer, notes and rubric

Design the Data Platform for Ride-Hailing Surge Pricing

Problem

Drivers send GPS pings every 4 seconds; riders open the app and request rides. The pricing service needs, per small geographic area, every ~30 seconds, the current supply (available drivers), demand (requests/app opens) and recent trend, to set a surge multiplier. Analysts need the history to tune the algorithm.


Clarifying questions

QuestionAssumed answer
Scale?5 M active drivers at peak globally, ping every 4 s; 30 M ride requests/day; 500 cities
Latency?Surge features updated every 30 s; end-to-end < 1 min
Area definition?Hexagonal cells (H3 resolution ~8, ~0.7 km²)
Consumers?Pricing service (online), driver heat map, analytics/ML
Late/out-of-order pings?Common (bad connectivity); ignore pings older than 2 min for real-time

1. Estimates

Pings: 5M / 4 s = 1.25M pings/s at peak × 100 B = 125 MB/s → Kafka 128–256 partitions
Requests/app-opens: ~5k/s
Active cells: ~1M worldwide; output 1M rows every 30 s → 33k rows/s to the feature store
Raw ping storage: 1.25M/s × 86,400 × 100 B ≈ 10 TB/day at peak rates → ~2 TB/day Parquet; retain raw 30–90 days, aggregates longer

2. Architecture

flowchart LR
    DRV[Driver app pings] --> GW[Edge gateway]
    RID[Rider app: opens, requests] --> GW
    GW --> KP[[Kafka: driver_pings<br/>key = driver_id]]
    GW --> KD[[Kafka: demand_events<br/>key = h3_cell]]
    KP --> F[Flink job]
    KD --> F
    F -->|"per cell, 30 s windows:<br/>supply, demand, ETA"| KO[[Kafka: cell_features]]
    KO --> ON[(Online store<br/>Redis / Cassandra<br/>key = city:cell)]
    ON --> PRICE[Pricing service]
    PRICE --> KS[[Kafka: surge_decisions]]
    KO --> PIN[(Pinot: heat maps,<br/>ops dashboards)]
    KP --> LAKE[(Delta bronze:<br/>pings, events, decisions)]
    KD --> LAKE
    KS --> LAKE
    LAKE --> ANA[Analytics, ML training,<br/>surge simulation / backtests]

3. Deep dives

3.1 Geospatial bucketing

3.2 Supply = state, not a count of pings

A driver pings 7–8 times per 30 s. Counting pings per cell would overcount, so supply is the number of distinct available drivers whose latest location is in the cell.

flowchart LR
    P[Ping: driver 42, cell A, status available] --> S{Keyed state<br/>by driver_id}
    S -->|"previous cell = B"| U1["supply[B] -= 1, supply[A] += 1"]
    S -->|"no ping for 60 s"| U2["timer fires: supply[A] -= 1<br/>(driver offline)"]
    U1 --> W[Emit per-cell snapshot every 30 s]
    U2 --> W

3.3 Event time and out-of-order pings

3.4 Feeding pricing safely

3.5 Analytics data model

4. Trade-offs

DecisionChoiceWhy
EngineFlinkKeyed state with timers (driver expiry), ms-level processing, high throughput
PartitioningPings by driver_id, aggregation by cellCorrect per-driver state, then re-key for per-cell sums
Online storeRedis (city-sharded)Low-latency reads by pricing; Cassandra if persistence + multi-DC needed
Spatial indexH3Uniform hexagons, hierarchical, open-source
Raw ping retention30–90 daysVolume is huge; keep derived intervals long-term

5. Failure modes

FailureHandling
Flink job restartState restored from checkpoint; supply counts consistent
City-wide network blip → no pingsTimers would mark all drivers offline → surge spikes! Detect sudden global supply drop, freeze surge / dampen changes
Hot cell (stadium event)Cell-level aggregation is cheap; skew is at driver level, which is uniform
Stale featuresPricing fallback with staleness threshold

6. What separates a senior answer

7. Follow-up questions

How would you compute "average ETA to pickup" per cell?

Join demand locations with nearby available drivers (k-ring) and road-network ETA estimates from a routing service; sample rather than compute for every request; aggregate per cell. Often a separate service owns ETA; the data platform provides features and logs.

How do you test a new surge algorithm without risking revenue?

Offline backtesting on logged features (counterfactual), then a switchback experiment (alternate algorithms by city × time slot, since user-level A/B testing causes marketplace interference), measuring conversion, wait times, driver earnings.


Self-assessment rubric