Skip to main content

Why fixed geohash grids create hot shards

This is the highest-signal follow-up in any geospatial design — Yelp, Uber, DoorDash, Tinder. The failure is specific and mechanical, and being able to walk through it beats naming H3 by a wide margin.

1. The mechanism​

How geohash works​

A geohash interleaves the bits of latitude and longitude and encodes the result in base 32. Each character adds 5 bits, so the string's length fixes the cell's physical size:

LengthApproximate cell sizeReal-world scale
439 km × 20 kmA metropolitan area
54.9 km × 4.9 kmA city district
61.2 km × 0.6 kmA few blocks
7153 m × 153 mA single large building
838 m × 19 mA building entrance

A natural design choice is: "partition the database by a 5-character geohash prefix — about 5 km cells."

Then compare two cells​

┌──────────────────────────────────┐ ┌──────────────────────────────────┐
│ "9q8yy" — rural Kansas │ │ "dr5ru" — Midtown Manhattan │
│ │ │ │
│ • 2 gas stations │ │ • 18,000 restaurants and bars │
│ • 1 diner │ │ • 45,000 active drivers │
│ │ │ • 150,000 concurrent searches │
│ Traffic: 0.1 QPS │ │ Traffic: 80,000 QPS │
│ Data: 10 KB │ │ Data: 1.5 GB │
└──────────────────────────────────┘ └──────────────────────────────────┘
Shard #1: idle Shard #2: on fire

Both cells cover the same physical area. Human activity within them differs by roughly six orders of magnitude. That ratio is the whole problem, and quoting it is far more persuasive than saying "density varies."

The three distinct failures​

Storage skew. One shard holds kilobytes, another holds gigabytes. Capacity planning becomes meaningless because the mean tells you nothing about the max.

Load skew — the one that actually takes you down. The Manhattan shard's CPU saturates, its connection pool exhausts, and its p99 climbs. Retries pile on. Because clients typically fan out to neighbouring cells, the failure spreads to adjacent shards, and a single geographic hot spot becomes a cluster-wide incident. Meanwhile 199 other nodes sit at 1% utilisation, so you cannot fix it by adding capacity — new nodes take a share of the keyspace, not a share of the load.

Query amplification if you shrink the cells. The obvious reaction is "use precision 7 so Manhattan splits into many small cells." Now consider a 10 km radius search in Kansas: at 153 m per cell, covering that circle requires thousands of cell lookups, nearly all returning nothing. You have traded a hot shard for a query that touches thousands of partitions. There is no single precision that serves both.

2. Why static quadtrees fail differently​

A quadtree recursively subdivides space into four quadrants, splitting only where density is high:

[ Root ]
┌────────┼────────┬────────┐
NW NE SW SE
┌─┼─┬─┐
... ... ... (Manhattan subdivides 15 levels deep)

This is genuinely adaptive — Kansas stays one large node, Manhattan subdivides to street level. So why is it not the answer?

Write-lock contention on splits. With moving entities (drivers, couriers, users), nodes must split and merge continuously. A split mutates the parent's child pointers and reallocates the children, requiring write locks up the tree. Under 50,000 location updates per second in one city, the tree spends its time contending on locks rather than serving queries.

It does not shard. A quadtree is a pointer-based structure. Distributing it across 50 machines means either replicating the whole tree (and then writes must be coordinated everywhere) or partitioning subtrees across nodes — at which point the Manhattan subtree lives on one machine and you are back to a hot shard, having added a distributed pointer structure for nothing.

Rebuilds are expensive and global. Recomputing the tree as the entity distribution shifts is a stop-the-world operation on a structure that is being read constantly.

3. The three production fixes​

Fix 1 — hierarchical multi-resolution indexing (H3 / S2)​

Do not pick one resolution. Index at a fine resolution and choose the query resolution per request, based on the search radius and local density.

import h3

def cells_for_query(lat, lng, radius_m):
"""Pick the coarsest resolution whose ring of cells covers the radius in a
modest number of cells. Dense areas naturally use finer resolutions
because their queries use smaller radii."""
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. In a square grid, the four edge-neighbours are at distance d and the four corner-neighbours at d√2 — so "the ring of adjacent cells" is not a circle, and distance filtering is uneven. Every neighbour of a hexagon is equidistant, which makes ring queries a much better approximation of a radius and keeps the candidate set tight.

H3 cell IDs also encode the resolution and parent chain in 64 bits, so coarsening is a bit operation rather than a lookup. S2 is the equally valid alternative: Hilbert-curve cells give excellent 1-D locality, which means spatial queries become range scans in any ordered store — a real advantage if your storage is SQL or a plain key-value database rather than a search engine.

Fix 2 — compound shard keys with density-proportional salting​

Hierarchy fixes query cost. It does not by itself fix load distribution, because a dense cell is still one key. Salt it:

Partition key = cell_id + "#" + salt

Rural Kansas: 9q8yy#0 (salt range 0..0 — one partition)
Manhattan: dr5ru#0, dr5ru#1, ... dr5ru#31 (salt range 0..31 — 32 partitions)
SALTS = load_salt_config() # per-cell, refreshed periodically from load metrics

def write_key(cell, entity_id):
return f"{cell}#{entity_id % SALTS.get(cell, 1)}"

async def read_cell(cell, filters):
n = SALTS.get(cell, 1)
parts = await asyncio.gather(*(store.query(f"{cell}#{i}", filters) for i in range(n)))
return heapq.nlargest(50, chain.from_iterable(parts), key=score) # scatter-gather

Two details that matter: the salt count must be stored, not inferred, because readers need to know how wide to fan out; and increasing a salt count is safe (fan out over the maximum historical value) while decreasing it requires a migration.

Fix 3 — sidestep spatial sharding entirely​

Often the best answer is to stop sharding by geography at all. Elasticsearch (and Lucene generally) hash-shards documents by ID, so load is uniform by construction regardless of geography, and every shard runs the geo predicate over its own slice using BKD-trees.

Geographic sharding: perfect locality, catastrophic skew, cheap queries
Hash sharding: zero locality, perfect balance, N-way fan-out per query

At 50M businesses across 20 shards, a 20-way fan-out per query is entirely affordable, and you have eliminated the hot-shard class of failure completely. Say which trade you are making and why — that is the answer, not the choice itself. Hash sharding wins when the corpus is small enough that fan-out is cheap; geographic sharding wins when the corpus is huge and queries must touch few partitions.

4. The boundary problem, while you are here​

Whatever indexing you choose, a user standing near a cell edge has their best result in the neighbouring cell. Always query a ring, and always re-filter by true distance:

cells = h3.grid_disk(origin_cell, k) # conservative SUPERSET of the circle
rows = await store.query_cells(cells, filters)
return sorted((haversine(user, r) , r) for r in rows if haversine(user, r) <= radius)[:50]

The direction of the approximation is the point: over-fetch a superset, then refine exactly. Under-fetching produces silently missing results, which nobody notices and nobody can debug.

5. What to say in an interview​

"A fixed-precision geohash as a shard key fails because human density spans about six orders of magnitude — a 5 km cell in Midtown holds 18,000 businesses and 80,000 QPS, the same-sized cell in rural Kansas holds three businesses and 0.1 QPS. The Manhattan shard saturates while 199 others idle, and adding nodes doesn't help because they take keyspace, not load. Shrinking the cells just converts that into query amplification: a 10 km rural search would scan thousands of empty cells.

Static quadtrees are adaptive in principle but fail on write-lock contention when entities move, and they don't shard — a city's subtree lands on one machine and you're hot again.

So: index hierarchically with H3 or S2 and pick the query resolution from the search radius; shard for uniform load rather than spatial locality — either hash-shard and fan out, or salt dense cells with a stored, density-proportional salt count; and always query a cell ring and re-filter by exact haversine, because over-fetching then refining is the only direction that can't silently lose results."