A CLICKHOUSE BATTLE STORY
Billboards near a
coffee shop
A battle story of making millions-scale proximity joins cheap in ClickHouse
Amaan Shaikh
Solutions Consultant at Sahaj Software
"Show me the best billboards within 1 Km of the starblobs coffee shops."
THE PROBLEM
1
2
3
4
5
THE PROBLEM
Billboards in NY within 1 Km of the coffee shop
Starblobs shop
1 Km
0.5 mile
Billboards
Places
1
2
3
4
5
New York City
1
2
3
4
5
THE PROBLEM
Billboards in NY & CA within 1 Km of the coffee shop
1
2
3
4
5
THE PROBLEM
Billboards in USA within 1 Km of the coffee shop
THE SOLUTION
First narrow down, then measure
Billion of pairs
A few thousands
The near ones
Cheap Filter
Exact Distance
1
2
3
4
5
Ways to group points by area
Method
What it does
The downside
Bounding box
a rectangle around the point
rough, has incorrect results
Geohash
rectangle cells named by a text code
rectangles, uneven at the edges
1
2
3
4
5
THE SOLUTION
THE SOLUTION
H3: Uber's open-source geospatial index.
1
2
3
4
5
THE SOLUTION
H3: Uber's open-source geospatial index.
1
2
3
4
5
Ways to group points by area
Method
What it does
The downside
Bounding box
a rectangle around the point
rough, has incorrect results
Geohash
rectangle cells named by a text code
rectangles, uneven at the edges
1
2
3
4
5
H3 grid
a hexagon as a plain number
boundary approximations
THE SOLUTION
Tag each point, join on the cell
billboard_id
lat
lon
h3
b_1042
40.71
-74.00
852a100bfffffff
b_1043
40.75
-73.98
852a1073fffffff
b_2984120
40.69
-73.91
852a100ffffffff
Billboards
place_id
lat
lon
h3
p_5501
40.72
-74.01
852a1051fffffff
p_5502
40.68
-73.95
852a1068fffffff
p_19847302
40.80
-73.90
852a100ffffffff
Places
SELECT b.billboard_id
FROM billboards b
JOIN places p ON b.h3 = p.h3
WHERE p.brand = 'starblobs'
AND geoDistance(b.lon, b.lat, p.lon, p.lat) <= 10000
THE SOLUTION
1
2
3
4
5
1
2
3
4
5
"Show me the best billboards around 1 Km of the starblobs coffee shops."
THE PROBLEM
1
2
3
4
5
What “best” means
The near billboards
Board
Reach
Cost
A
8
$$$
B
4
$
C
9
$$$$
D
7
$$
E
6
$$
many signals per billboard
Score,
Then sort
Weights change
every search
Goal: reach
1
Billboard C
2
Billboard A
3
Billboard D
4
Billboard E
5
Billboard B
best at the top
Goal: budget
1
Billboard B
2
Billboard D
3
Billboard E
4
Billboard A
5
Billboard C
best at the top
1
2
3
4
5
Existing End to End Flow
User Interface
Proximity service
Reads Redis, computes distance
Ranking service
Scores the billboards
Mongo
Places (brand -> lat, lon)
Redis
Billboard geo index lookup
THE SOLUTION
1
2
3
4
5
ClickHouse
Billboards and ranks
End to end query benchmarks
1
2
3
4
5
1
2
3
4
5
1
2
3
4
5
WHY IT BROKE
1
2
3
4
5
1
Slow queries
2
End-to-end is far worse than the query
Spike to compare all the contenders
1
2
3
4
5
1
2
3
4
5
1
2
3
4
5
THE REBUILD
Why is ClickHouse so fast?
billboard_id
b_1
b_2
b_3
b_4
b_5
reach
31
27
44
22
38
h3
85f05ab7fffff
Only the column you need
One sweep
THE REBUILD
1
2
3
4
5
1
2
3
4
5
1
2
3
4
5
85f05ab7fffff
85f05ab7fffff
85f05ab7fffff
85f05ab7fffff
SIMD: one instruction, many values
4
7
2
9
1
8
3
5
One CPU instruction -> all 8 at once
Place
k = 0 the place's cell
k = 1 the first ring
k = 2 a wider ring
End to end query benchmarks with ClickHouse
1
2
3
4
5
1
2
3
4
5
1
2
3
4
5
THE REBUILD
New End to End Flow
User Interface
Ranking service
Scores the billboards
ClickHouse
Billboards, Places and ranks
1
2
3
4
5
1
2
3
4
5
1
2
3
4
5
THE REBUILD
THE TAKEAWAYS
1
Know when and how much to precompute.
2
Proximity is an equality join on the cell, not a distance filter.
3
Do the join in a columnar engine, next to the data.
1
2
3
4
5
Questions?
Amaan Shaikh
Solutions Consultant at Sahaj Software
THE TAKEAWAYS
Change the condition,
not the query.
A distance range became an equality on a cell, and the work we could never skip simply went away.
Know when and
how much to precompute.
Precomputing the answer was fast per query but huge and rigid. Weigh every constraint, not just speed.
1
2
3
4
5
Two problems with the current solution
1
Slow queries
A single proximity query already runs into seconds. The hybrid recomputes
distance across millions of candidate pairs, one row at a time.
2
End-to-end is far worse than the query
Shipping 100-200k candidate ids across services and regions, then re-ranking,
pushes the round-trip to minutes. The widest queries are not supported at all.
WHY IT BROKE
1
2
3
4
5
1
2
3
4
5
1
2
3
4
5
1
2
3
4
5
ClickHouse wins the queries people actually run
Scenario
Before
Mongo
Redis
ClickHouse
Small brand, 1 mi
3 s
2.6 s
1.9 s
0.08 s
Small brand, 10 mi
3.5 s
3.9 s
1.6 s
0.11 s
Mid brand, 1 mi
4 s
8.4 s
5.7 s
0.26 s
Mid brand, 10 mi
4.5 s
8.7 s
4.2 s
0.89 s
Large brand, 1 mi
14.9 s
18.1 s
19.9 s
4.4 s
Dense brand, 10 mi
25.5 s
12.4 s
7.0 s
5.5 s
Fastest in every realistic query. The lone exception: the densest category at a wide radius.
THE REBUILD
Agenda
1
The problem
2
The solution
3
Why it broke?
4
The rebuild
5
The takeaways
THE PROBLEM
Billboards and places
What we do
Plan out-of-home campaigns which billboards and
screens a brand should book.
Billboards and places
Runs on two sets of points on the map: the billboards
we can book, and the places people go to.
Billboards
Places
0.5 mile
1
2
3
4
5
New York City
Our coffee shop
0.5 mile
1
2
3
4
5
THE PROBLEM
Billboards in NY around 2 Km of the coffee shop
New York City
2 Km
Billboards
Places
ClickHouse: granules + SIMD
Column on disk
Index
Skip
Skip
Read
Skip
Skip
Granule = 8,192 rows
1 mark each
Read one granule
SIMD: one instruction, many values
4
7
2
9
1
8
3
5
One CPU instruction -> all 8 at once
8
14
4
18
2
16
6
10
A scalar CPU would do these one at a time
THE REBUILD
1
2
3
4
5
1
2
3
4
5
1
2
3
4
5
THE SOLUTION
Two datasets
billboard_id
lat
lon
b_1042
40.71
-74.00
b_1043
40.75
-73.98
b_1044
40.73
-73.99
b_1045
40.72
-74.02
…
…
…
b_2984120
40.69
-73.91
Billboards
≈ 3 million rows
place_id
lat
lon
p_5501
40.72
-74.01
p_5502
40.68
-73.95
p_5503
40.71
-73.96
p_5504
40.66
-73.93
…
…
…
p_19847302
40.80
-73.90
Places
≈ 20 million rows
1
2
3
4
5
How the old system answered it
Proximity service
Reads Redis, computes distance
1. Near StarBlobs radius R
4. 100-200k billboard-ids
Mongo
Places (brand -> lat, lon)
Redis
Billboard geo index lookup
THE SOLUTION
1
2
3
4
5
2. Get Places
3. Get billboards matching h3
Why not Redis or Mongo?
The proximity query is the hard part, and neither datastore is specifically built for it.
REDIS
key-value store
A bucket lookup: ids in a cell, nothing more.
No arithmetic: cannot compute the distance.
Cannot filter candidates by the actual radius.
Finds candidates. Cannot measure the distance.
MONGODB
document store
Computes distance one document at a time
Per-document cost across millions of pairs.
No columnar layout, no vectorized distance.
Computes distance, but far too slowly at scale.
THE REBUILD
1
2
3
4
5
1
2
3
4
5
1
2
3
4
5
Query benchmarks for Proximity service
THE SOLUTION
1
2
3
4
5
THE SOLUTION
1
2
3
4
5
res
h3
cell area
0
~4,300,000 sq km
…
…
…
4
~1,800 sq km
5
~250 sq km
6
~36 sq km
…
…
…
15
~0.9 sq m
802a
fffffffffff
842a107
ffffffff
852a1073
fffffff
862a10737
ffffff
8f2a1073b59c2d4
1
2
3
4
5
H3: Uber's open-source geospatial index.
THE ATTEMPTS
Precompute the pairs
Precompute every billboard-place pair within ten miles, once.
We did the napkin math before writing any code.
~8 hours
To build the table
~600 GB
Of pairs to store
Costly
To serve it fast
And the radius is baked in, so a new radius rebuilds all of it.
Too big, too rigid, not cost-effective. We never built it.
1
2
3
4
5
THE ATTEMPTS
Just run the query
Skip precompute. Just measure the distances at query time, with a plain join.
One nationwide query, for a big coffee-shop brand:
35,000 coffee shops
one brand, nationwide
×
3 million billboards
the whole country
=
105 billion
Distance checks, one query
No index skips a distance you have to compute. Too slow, every time.
1
2
3
4
5
THE ATTEMPTS
The real problem
Both attempts were built on the same join condition.
ON distance(billboard, place) <= radius
A range
Precompute
Ran it for every pair, up front
Too costly
Just run it
Ran it for every pair, live
Too slow
A range is something you can't index or skip.
1
2
3
4
5
THE CONSTRAINTS
What we were up against
Anywhere from 1 to 10 miles
The search radius is
never the same
Fast under load
Subsecond to seconds,
20 to 50 users
Cheap to run
No big new machines
1
2
3
4
5
THE SOLUTION
What’s H3
H3 is Uber's open-source geospatial index.
1
2
3
4
5
geoToH3
(40.71, -74.00, 5)
=
852a1073fffffff
res
h3
cell area
0
~4,300,000 sq km
…
…
…
4
~1,800 sq km
5
~250 sq km
6
~36 sq km
…
…
…
15
~0.9 sq m
802a
fffffffffff
842a107
ffffffff
852a1073
fffffff
862a10737
ffffff
8f2a1073b59c2d4
1
2
3
4
5
THE SOLUTION
The resolution is a knob
geoToH3
(40.71, -74.00, 5)
=
852a1073fffffff
res
h3 cell
cell area
0
~4,300,000 sq km
…
…
…
4
~1,800 sq km
5
~250 sq km
6
~36 sq km
…
…
…
15
~0.9 sq m
802a
fffffffffff
842a107
ffffffff
852a1073
fffffff
862a10737
ffffff
8f2a1073b59c2d4
1
2
3
4
5
Query benchmarks
Scenario
Places matched
Query time without cache
Small brand, 1 mi
640
3 s
Small brand, 10 mi
640
3.5 s
Mid brand, 1 mi
4.2K
4.5 s
Mid brand, 10 mi
4.2K
7 s
Large brand, 1 mi
66K
14.9 s
Dense brand, 10 mi
170K
25.5 s
THE SOLUTION
1
2
3
4
5
THE SOLUTION
One cell is never enough
1
2
3
4
5
The edge effect
The place's cell
Closer, but missed
The cell and its ring
Now caught
Add the ring
THE SOLUTION
What’s H3
H3 is Uber's open-source geospatial index.
1
2
3
4
5
THE SOLUTION
The ring is a knob too, in ClickHouse
The ring we just added is one call: h3kRing.
It takes a cell and k, and returns the cell
plus k rings.
Dial resolution and k to fit the radius. At
res 5, k = 2 is about 10 miles.
SELECT b.billboard_id
FROM billboards b
JOIN places p ON b.h3 IN h3kRing(p.h3, 2)
WHERE p.brand = 'coffee shops'
AND geoDistance(b.lon, b.lat, p.lon, p.lat) <= 16093
Place
k = 0 the place's cell
k = 1 the first ring
k = 2 a wider ring
1
2
3
4
5
THE REBUILD
1
2
3
4
5
1
2
3
4
5
1
2
3
4
5
CLOAKROOM
#1
#2
#3
#4
#5
#6
#7
#8
#9
Your token
#5
Present it
Get bag #5
Your exact bag
Your app
Changing it = compute
Redis, like a cloakroom
THE REBUILD
1
2
3
4
5
1
2
3
4
5
1
2
3
4
5
MongoDB, like folders in a cabinet
MONGODB
user:11
user:53
user:88
user:07
user:42
user:90
user:19
user:64
user:97
Find by key
The whole document
{
"id": 42,
"name": "Alex",
"city": "Pune",
"orders": [ 3 items ],
"vip": true
}
THE SOLUTION
Attempts vs the solution
distance(billboard, place) <= radius
billboard.h3 = place.h3
600 GB, or 105 billion checks a query
Radius baked in
Minutes, or no answer
~1 GB of points
Any radius, just dial k
Answers in seconds
Change the condition, and everything downstream shrinks.
1
2
3
4
5
THE SOLUTION
Coffee shops, nationwide
The widest query there is: every billboard near any coffee shop, nationwide.
Within 10 miles
10-15s
On ClickHouse
1
2
3
4
5
THE SOLUTION
End to end
Precompute
Tag every billboard and place
with its H3 cell
Store
Store the points and their
cells. ~1GB, not pairs
Serve
Expand the k-ring, join on
cell, exact distance on the
few
Points, not pairs. About 1GB, not 600GB.
1
2
3
4
5