Skip to main content

Yelp / Nearby Places (Proximity Service)

Sharpened prompt. Design a service that answers "restaurants within 2 km, rating above 4, price ≤ $$, open right now, ranked for me" in under 100ms, over 50M businesses, at 100k QPS — where Manhattan has 18,000 candidates in a square kilometre and rural Kansas has three in a hundred.

The whole problem is that humans are not uniformly distributed but grids are. Every naive spatial partitioning scheme fails on exactly that mismatch.

1. Problem framing​

Functional requirements​

  • Search businesses by location and radius (or map viewport).
  • Filter by category, rating, price, hours, attributes (outdoor seating, delivery).
  • Rank by relevance: distance, rating, popularity, personalisation.
  • Write path: new businesses, reviews, photos, hour changes.

Non-functional requirements​

PropertyTargetConsequence
Query latencyP99 < 100msIndex must live in memory, filters must be indexed
Read:write~1000:1Immutable, read-optimised index structures are the right call
FreshnessMinutes for edits, seconds for "open now"Time-dependent filters must be computed, not indexed as booleans
Correctness at boundariesNo missing results near cell edgesNeighbour cells must be part of every query

Back-of-the-envelope​

Businesses: 50M globally, ~500 B of indexable attributes = 25 GB -> fits in RAM
Reviews: 500M × 1 KB = 500 GB (separate store, not in the search index)
Queries: 100k QPS, 90% within the 10 densest metro areas
Density: Manhattan ~18,000 businesses/km²
Rural KS ~0.03 businesses/km²
That is a 600,000× ratio. No fixed cell size works for both.

Quote the density ratio. It is the single most persuasive way to show why the obvious answer fails.

2. High-level architecture​

3. Component inventory​

ComponentConcrete choiceWhy this one
Spatial indexUber H3 (hexagons) or Google S2 (Hilbert-curve cells)Hierarchical, multi-resolution, and H3's equidistant neighbours make ring queries clean
Search engineElasticsearch / LuceneBKD-trees for geo and numeric ranges, inverted indexes for attributes, combined in one query
Shard keyH3 cell at a coarse resolution, salted in dense cellsSee 4.2 — this is where the hot-shard fix lives
Result cacheRedis, keyed by (cell_set, filter_hash, sort)Map viewports repeat heavily; a snapped viewport is cacheable
Write pipelineKafka → segment builder → blue/green swapImmutable segments make rebuilds safe and atomic
RankingLightweight in-shard, personalised at the mergerPersonalisation cannot be baked into a shared index

4. The toughest parts​

4.1 Why fixed grids fail — and what replaces them​

Why it's hard. Geohash encodes a lat/long into a base-32 string; the string's length fixes the cell size. Pick length 5 (≈5 km × 5 km) as a shard key and you get:

┌──────────────────────────────┐ ┌──────────────────────────────┐
│ "9q8yy" — rural Kansas │ │ "dr5ru" — Midtown Manhattan │
│ 3 businesses │ │ 18,000 businesses │
│ 0.1 QPS │ │ 80,000 QPS │
│ 10 KB │ │ 1.5 GB │
└──────────────────────────────┘ └──────────────────────────────┘
Shard 1: idle Shard 2: on fire

Shrink the cells to fix Manhattan and rural queries must scan thousands of empty cells to cover a 10 km radius — query amplification replaces storage skew. There is no single resolution that works, because the underlying distribution spans six orders of magnitude.

Static quadtrees are adaptive in principle (subdivide only where dense) but fail differently: they are pointer structures, so distributing one across 50 machines means a city's subtree lives on one node (hot again), and dynamic splits require write locks across parent and child nodes, which under 50k location writes/sec is permanent lock contention.

Solution — a hierarchical multi-resolution index, queried at a resolution chosen per request.

H3 indexes the globe with hexagons at 16 resolutions. Each cell has a 64-bit ID that encodes its resolution and its parent chain, so coarsening is a bit shift rather than a lookup.

import h3

# Index each business at a fine resolution; the hierarchy gives you the rest for free.
def index_business(b):
b.h3_12 = h3.latlng_to_cell(b.lat, b.lng, 12) # ~9 m edge, the storage key
# Parents are derivable, so we store a few for cheap coarse filtering.
b.h3_9 = h3.cell_to_parent(b.h3_12, 9) # ~174 m
b.h3_7 = h3.cell_to_parent(b.h3_12, 7) # ~1.2 km

# Query: pick the resolution so the ring covers the radius in a modest number of cells.
def cells_for_radius(lat, lng, radius_m):
for res in range(12, 4, -1):
edge = h3.average_hexagon_edge_length(res, unit='m')
k = math.ceil(radius_m / edge)
if k <= 3: # 1 + 3k(k+1) <= 37 cells
origin = h3.latlng_to_cell(lat, lng, res)
return h3.grid_disk(origin, k), res
return [h3.latlng_to_cell(lat, lng, 5)], 5

Why hexagons rather than squares: every neighbour of a hexagon shares an edge and sits at the same centre-to-centre distance. With squares, diagonal neighbours are √2 farther, so "ring of k cells" is not a circle and distance filtering is uneven. For a proximity service, that uniformity directly improves both correctness and query cost.

S2 is the equally valid alternative — Hilbert-curve cells give excellent locality and native range-query support in any ordered store (WHERE s2_cell BETWEEN lo AND hi), which is a real advantage if your storage is a plain key-value or SQL database rather than a search engine.

4.2 Sharding a skewed world without hot shards​

Why it's hard. Even with a good spatial index, if the shard key is the cell ID, Manhattan's cell is one shard and it takes a disproportionate share of both data and traffic. Hierarchy alone does not fix distribution.

Solution — compound keys with density-proportional salting, plus scatter-gather reads.

# Salt count is proportional to observed density, recomputed periodically.
SALTS = {"dr5ru": 32, "9q8yy": 1} # Manhattan: 32 partitions; Kansas: 1

def write_partition(cell: str, business_id: int) -> str:
n = SALTS.get(cell, 1)
return f"{cell}#{business_id % n}" # deterministic, so reads know the range

async def read_cell(cell: str, filters) -> list:
n = SALTS.get(cell, 1)
# Scatter-gather across the salt range, merge with a bounded heap.
parts = await asyncio.gather(*(
store.query(f"{cell}#{i}", filters) for i in range(n)))
return heapq.nlargest(50, chain.from_iterable(parts), key=lambda b: b.score)

The salt count must be stored, not inferred, so readers know exactly how many partitions to fan out to. Keep it in a small config table refreshed every few minutes; increasing a salt count is a safe operation (old data stays reachable if you fan out over the max historical value), decreasing it requires a migration.

In practice, if you use Elasticsearch you get most of this for free: it document-shards by hash of the document ID, so businesses distribute uniformly regardless of geography, and the geo query runs on every shard. You trade fan-out cost for guaranteed uniform load — usually the right trade at 50M documents, which is small enough that a 20-shard fan-out is cheap. Be explicit about which trade you are making.

4.3 The boundary problem​

Why it's hard. A user stands 50 metres from the edge of their cell. The best restaurant is 100 metres away, in the neighbouring cell. Query only the containing cell and you miss it — and the failure is invisible, because you return some results and nobody knows what is missing.

Solution — always query a ring, and always re-filter by true distance.

async def nearby(lat, lng, radius_m, filters):
cells, res = cells_for_radius(lat, lng, radius_m) # grid_disk already covers the ring
rows = await store.query_cells(cells, filters)

# The cell set is a superset of the circle. Compute exact distance and cut.
out = []
for r in rows:
d = haversine(lat, lng, r.lat, r.lng)
if d <= radius_m:
out.append((d, r))
return sorted(out)[:50]

Two properties make this correct and fast: the cell ring is a conservative superset of the circle (never fewer cells than needed, so no false negatives), and the exact distance filter is applied to a small candidate set. Getting the direction of the approximation right — over-fetch then refine, never under-fetch — is the thing to say.

For map viewport queries (a rectangle, not a circle), use h3.polygon_to_cells to cover the polygon, and snap the viewport to a coarse grid before caching so that slight pans reuse the same cache key rather than generating a new one every frame.

4.4 Geo plus five other filters, at once​

Why it's hard. "Within 2 km" is one predicate; "rating > 4, price ≤ $$, category = Thai, open now, has outdoor seating" are five more. Filtering geo first then scanning attributes means scanning 18,000 Manhattan candidates per query. Filtering attributes first means scanning every Thai restaurant on Earth. Neither ordering is universally right — it depends on selectivity, which varies per query.

Solution — put everything in one engine that can choose the intersection order by cost.

{
"query": {
"bool": {
"filter": [
{"terms": {"h3_9": ["8928308280fffff", "..."]}},
{"term": {"category": "thai"}},
{"range": {"rating": {"gte": 4.0}}},
{"range": {"price_tier": {"lte": 2}}},
{"term": {"attributes": "outdoor_seating"}},
{"geo_distance": {"distance": "2km", "location": {"lat": 40.75, "lon": -73.98}}}
]
}
},
"sort": [{"_score": "desc"}],
"size": 50
}

Lucene handles this well because it has the right structures for each predicate: BKD-trees for the numeric and geo ranges, inverted postings for the terms, and a cost-based planner that intersects the most selective clause first. Providing the H3 cell terms alongside geo_distance gives the planner a cheap, highly selective pre-filter — the cell terms narrow the candidate set with a term lookup, and the exact distance check runs on the survivors.

"Open now" cannot be indexed as a boolean, because it changes every minute and would require re-indexing 50M documents continuously. Index the schedule as numeric ranges instead (minutes-since-Monday-midnight intervals) and let the query supply the current time:

{"bool": {"should": [
{"bool": {"filter": [
{"range": {"hours.open_min": {"lte": 1043}}},
{"range": {"hours.close_min": {"gt": 1043}}}
]}}
]}}

Same trick for "delivers to me" and any other time- or viewer-dependent predicate: index the invariant, evaluate the variable at query time.

4.5 Rebuilding the index without downtime​

Why it's hard. Categories get reorganised, the ranking model changes, a new attribute is added, or the H3 resolution strategy is tuned. Any of these requires a full re-index of 50M documents. Doing it in place means the index is inconsistent mid-flight — some documents in the new format, some in the old — and mutating a live index while serving 100k QPS causes latency spikes from merge pressure.

Solution — immutable index versions with an atomic pointer swap.

1. Build v_N+1 in the background from the source of truth, at whatever pace is cheap.
2. Tail the Kafka mutation stream into v_N+1 from the build's start offset, so it
catches up to live while v_N still serves.
3. Shadow-query: mirror 1% of live traffic to v_N+1, compare result sets, alert on
recall or ranking deltas beyond a threshold.
4. Atomic swap: repoint the alias. Readers pick it up on their next query.
5. Keep v_N warm for one hour. Rollback is another pointer swap.
# Elasticsearch makes the swap a single atomic action.
curl -XPOST localhost:9200/_aliases -d '{
"actions": [
{"remove": {"index": "biz_v41", "alias": "biz"}},
{"add": {"index": "biz_v42", "alias": "biz"}}
]
}'

The shadow-query step (3) is the one candidates omit and the one that saves you: an index rebuild that silently changes recall is very hard to detect from aggregate latency and error metrics, because nothing errors — results are just quietly worse. Compare result sets, not just health checks.

4.6 Ranking, personalisation, and the review economy​

Why it's hard. Distance and rating alone produce bad results: a 5.0-rated place with three reviews should not outrank a 4.6 with 2,000, and a spammy new listing should not outrank an institution. Meanwhile ratings themselves are under adversarial pressure — businesses buy reviews, competitors leave fake negatives.

Solution — a confidence-adjusted quality score, computed offline, plus personalisation at merge time.

def quality(biz) -> float:
# Wilson lower bound: shrinks small samples toward the mean. A 5.0 with 3 reviews
# scores below a 4.6 with 2,000, which is the correct ordering.
n, p = biz.review_count, biz.positive_ratio
if n == 0: return 0.0
z = 1.96
denom = 1 + z*z/n
centre = p + z*z/(2*n)
margin = z * math.sqrt((p*(1-p) + z*z/(4*n)) / n)
return (centre - margin) / denom

def rank(biz, viewer, dist_m) -> float:
return (0.35 * quality(biz)
+ 0.30 * distance_decay(dist_m) # exp(-d/scale), not linear
+ 0.15 * popularity_prior(biz) # views, saves, directions
+ 0.20 * viewer_affinity(viewer, biz)) # cuisine history, price band

distance_decay should be exponential, not linear: the difference between 200 m and 400 m matters much more than between 5 km and 5.2 km.

Review integrity is its own pipeline and worth a sentence: graph clustering over reviewer–business bipartite edges (rings of accounts reviewing the same set of businesses), burst detection (20 five-star reviews in an hour after months of silence), and text similarity across reviews. Suspect reviews are excluded from the score but remain visible, which avoids a false-positive removal turning into a customer-support incident.

Personalisation must happen after the shard merge, since it is viewer-specific and cannot be baked into a shared index — the same constraint as post search.

4.7 Caching queries that are almost, but not quite, identical​

Why it's hard. Two users standing ten metres apart issue different queries (different lat/long), so a naive cache key never hits. But their result sets are essentially identical. Meanwhile a map application fires a new query on every pan and zoom frame.

Solution — snap the query to a canonical form before caching.

def cache_key(lat, lng, radius_m, filters, sort) -> str:
# Snap the origin to a coarse H3 cell: everyone in the same ~170 m hexagon
# shares a cache entry. Snap the radius to a fixed ladder.
cell = h3.latlng_to_cell(lat, lng, 9)
r = min(x for x in (500, 1000, 2000, 5000, 10000) if x >= radius_m)
return f"nearby:{cell}:{r}:{stable_hash(filters)}:{sort}"

Snapping trades a little precision for an enormous hit-rate gain — and the precision loss is recovered because you still compute exact distances during hydration and re-sort the cached candidate IDs by the user's true position. Cache the candidate ID list, not the ordered rendered results, so the snap does not degrade the final ordering.

Time-dependent filters (open_now) must be part of the key or excluded from the cached set; the clean approach is to cache the geo+attribute candidates and apply the time predicate at hydration, keeping the cached object time-invariant.

5. What breaks first​

EventFirst failureMitigation
Dense metro at peakOne shard's CPU (if geo-sharded)Hash sharding or density-proportional salting
Map pan/zoom floodQuery rate multiplies 10× per sessionViewport snapping + client-side debounce + candidate caching
Index rebuildMerge pressure spikes latencyBlue/green with alias swap; never rebuild in place
A city-wide event (marathon)Cache miss storm in one areaPre-warm the affected cells; singleflight on cache fill
Review bomb on one businessRanking distortion, moderation queueWilson score damps small bursts; burst detection quarantines
Deep pagination on a sparse areaFan-out scans many empty cellsCoarsen resolution adaptively; cap page depth

6. Cheat sheet​

  • The number: Manhattan vs rural Kansas is a 600,000× density ratio. No fixed cell size survives it.
  • Index: H3 (or S2) hierarchical cells; index fine, query at a resolution chosen from the radius.
  • Sharding: hash-shard for uniform load, or salt dense cells proportionally with a stored salt count.
  • Boundaries: always query the ring (grid_disk), always re-filter by exact haversine. Over-fetch then refine.
  • Filters: one engine (Lucene) with BKD + inverted postings; index the schedule, evaluate "open now" at query time.
  • Rebuilds: immutable versions, shadow-query comparison, atomic alias swap, keep the old version warm.
  • Caching: snap the origin to a coarse cell and the radius to a ladder; cache candidate IDs, re-sort exactly.
  • The one-liner: "The Earth is not uniform, so no single-resolution grid works — I index hierarchically, choose resolution per query, shard for uniform load rather than spatial locality, and always over-fetch a cell ring before exact filtering."