skip to content

In a ride-hailing app that shards its in-memory driver location index by region, what problems do shard boundaries cause and how are they handled?

level: seniorimportance: nice to knowfreq 30%

answer

  1. three edge effects
  2. write new before deleting old
  3. expiry as the cleanup backstop
  4. search circle overlaps the neighbour
  5. a margin before switching shards

basics

~20 s

Moving drivers leave stale entries in the old shard, searches near an edge miss drivers across it, and drivers on a border flip between shards. Fixes are handoff plus dedupe, fan-out to every overlapping shard, and hysteresis.

solid answer

~40 s

Each shard owns an area (a city, or a group of cells), and a router maps a position to its shard. Boundaries cause three problems. **Handoff**: after a crossing, the new shard gets the update but the old one still holds an entry, so the router writes to the new shard, then deletes from the old one, with the old entry's time-to-live as a backstop; merges **dedupe by driver id**, keeping the newest timestamp. **Edge searches**: a search near a border crosses into the neighbour, so the dispatcher **fans out** to every shard the search area overlaps and merges. **Flapping**: a driver on a border road switches shard on every ping, so **hysteresis** keeps them in the current shard until they are clearly past the edge. Drawing borders through sparse areas reduces all three.

go deeper

for a junior

Recall that splitting a map into regions creates edges, and that a nearby driver can sit just across one.

for a middle

Explain the three edge effects - handoff, edge searches, flapping - and name one fix for each.

for a senior

Show the ordering and backstops that make handoff safe, how fan-out is bounded with timeouts, and why dispatch state stays outside the regional index.

for a principal

Discuss where to draw borders, what places force constant fan-out, and why a rebuildable index makes eventually consistent boundary handling acceptable.

## Why shard the live index by region A **ride-hailing app** keeps driver positions in an **in-memory geo index**. One index for a whole country would need to fit in one server's memory and absorb every update. So the index is **sharded by region**: each shard owns an area - a metro region, or a contiguous group of cells - and holds only the drivers inside it. A **router** maps any coordinate to its owning shard. This works well because nearly every query is local: a pickup only cares about drivers within a few kilometres. The trouble is concentrated at the **edges**. (Splitting an overloaded shard is a general sharding topic and is not covered here.) ## Problem 1: handoff when a driver crosses A driver in shard A drives across the border. The next update is routed to shard B, which inserts the driver. Shard A still holds the previous entry. - For up to one time-to-live, the driver exists **twice**: a stale copy in A and a live copy in B. - A search near the border may return both copies. - A search deep inside A may return a **ghost** that is really somewhere else. Handling: 1. The router remembers each driver's **current shard** (in the session for a persistent connection, or in a small lookup). 2. When the computed shard changes, it writes to the new shard **first**, then sends a delete to the old shard. Writing first means the driver is never missing, only briefly duplicated. 3. If the delete is lost, the old entry's **time-to-live** removes it anyway. 4. Anything that merges results **dedupes by driver id** and keeps the entry with the newest timestamp. ## Problem 2: searches that straddle an edge A pickup sits 300 m inside shard A, and the dispatcher searches 2 km around it. Almost half of that circle lies in shard B. Querying only A misses every driver across the border, even ones 400 m away. Handling: - compute the **set of shards whose area intersects the search area** (a bounding box or circle against each shard's area); - **fan out** the query to each of them in parallel; - **merge**, dedupe by driver id, and keep the best K; - set a per-shard timeout so one slow neighbour degrades the result instead of blocking it. Most searches touch one shard; only the ones near an edge pay the fan-out cost. ## Problem 3: flapping on the border A driver on a road that follows the border may be computed into A, then B, then A on successive pings. Each flip causes a handoff: a write, a delete and a brief duplicate. Handling: **hysteresis**. A driver stays in the current shard until they are more than some margin, say 200 m, past the border. The cost is that a shard can hold drivers slightly outside its area, so edge detection for fan-out must use each shard's area **expanded by that margin**. ## Summary | Problem | Symptom | Mitigation | |---|---|---| | Handoff | Duplicate or ghost driver after a crossing | Write new shard, then delete old; time-to-live backstop; dedupe by id | | Edge search | Missed drivers just across the border | Fan out to all overlapping shards and merge | | Flapping | Constant handoffs on a border road | Hysteresis margin; expand areas by it for fan-out | ## Choosing where borders sit Borders are cheapest where little happens: - follow **natural gaps** - rivers, parks, sparsely populated land - rather than cutting through a dense downtown; - keep a **metro area in one shard** when it fits, since most trips start and end inside it; - be careful with places that draw traffic from two regions, such as an airport between two cities, because every search there fans out. ## What must not depend on the shard Dispatch state - offers and assignments - should live in its authoritative store, **not** in the regional index. Then a handoff can only duplicate or briefly misplace a **position**, which the next ping repairs; it can never lose or duplicate an **assignment**. Keeping the index purely rebuildable is what makes these simple, eventually consistent boundary fixes acceptable.

  • Why write to the new shard before deleting from the old one during a handoff?
    If the delete happened first and the write then failed or lagged, the driver would be missing from every shard and would receive no offers until the next ping. Writing first means the worst case is a brief duplicate, which merges already handle by deduping on driver id and keeping the newest timestamp.
  • How does hysteresis change the fan-out rule for edge searches?
    With a margin, a shard may hold drivers up to that distance outside its own area. A search must therefore compare its area against each shard's area expanded by the margin; otherwise a driver still held by the neighbouring shard, but physically inside this shard's area, would be missed.

saying these in an interview costs you the question

  • A search only needs to query the shard that contains the pickup point.
  • Deleting from the old shard before writing to the new one is the safer order.
  • A driver can never appear in two shards at the same time.
  • Shorter update intervals fix drivers flapping between shards.
  • Assignments can safely live inside the regional location index.