1 of 48

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

2 of 48

"Show me the best billboards within 1 Km of the starblobs coffee shops."

THE PROBLEM

1

2

3

4

5

3 of 48

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

4 of 48

1

2

3

4

5

THE PROBLEM

Billboards in NY & CA within 1 Km of the coffee shop

5 of 48

1

2

3

4

5

THE PROBLEM

Billboards in USA within 1 Km of the coffee shop

6 of 48

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

7 of 48

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

8 of 48

THE SOLUTION

H3: Uber's open-source geospatial index.

1

2

3

4

5

9 of 48

THE SOLUTION

H3: Uber's open-source geospatial index.

1

2

3

4

5

10 of 48

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

11 of 48

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

12 of 48

"Show me the best billboards around 1 Km of the starblobs coffee shops."

THE PROBLEM

1

2

3

4

5

13 of 48

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

14 of 48

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

15 of 48

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

16 of 48

Spike to compare all the contenders

1

2

3

4

5

1

2

3

4

5

1

2

3

4

5

THE REBUILD

17 of 48

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

18 of 48

End to end query benchmarks with ClickHouse

1

2

3

4

5

1

2

3

4

5

1

2

3

4

5

THE REBUILD

19 of 48

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

20 of 48

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

21 of 48

Questions?

Amaan Shaikh

Solutions Consultant at Sahaj Software

22 of 48

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

23 of 48

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

24 of 48

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

25 of 48

Agenda

1

The problem

2

The solution

3

Why it broke?

4

The rebuild

5

The takeaways

26 of 48

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

27 of 48

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

28 of 48

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

29 of 48

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

30 of 48

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

31 of 48

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

32 of 48

Query benchmarks for Proximity service

THE SOLUTION

1

2

3

4

5

33 of 48

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.

34 of 48

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

35 of 48

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

36 of 48

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

37 of 48

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

38 of 48

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

39 of 48

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

40 of 48

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

41 of 48

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

42 of 48

THE SOLUTION

What’s H3

H3 is Uber's open-source geospatial index.

1

2

3

4

5

43 of 48

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

44 of 48

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

45 of 48

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

}

46 of 48

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

47 of 48

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

48 of 48

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