← all field notes
essay / October 10, 2026•26 min read

We built log search on ClickHouse instead of running Elasticsearch. Here's what we learned.

clickhouseloggingobservabilitydatabasesparsers

One field, two types, zero results

It started like most bad mornings do: with a search that returned nothing.

↳ the page

  • 10:42A deploy goes out. Nothing looks wrong.
  • 11:15Someone searches the logs for a failing order. Nothing comes back, even though the service clearly logged it.
  • 11:30The log shipper's own logs are full of mapping errors.
  • 11:40The cause turns up. One service logs "status": 200. A newer one logs "status": "OK". The index had already decided status is a number, so every line with the string version was rejected.

The logs existed. The service wrote them, the shipper read them, and the store refused to keep them, because two teams disagreed about the type of one field. Nobody was paged for the rejected lines. They were just gone, and we found out only because someone went looking for one.

That is the problem this post is about. Logs have no fixed format. Some are plain strings, some are JSON, some are key=value, some are web server access lines from 1995. The same field name means different things in different services, and new fields appear with every deploy. Any store that wants a schema up front will, sooner or later, reject or mangle real logs.

We built a log store where a type disagreement is not an error, on top of ClickHouse, a columnar database. This post explains the ideas that made it work, the ones that surprised us, and what we gave up. It does not publish our schema or our query language. Every code snippet is a simplified illustration, labeled that way, and there are seven interactive visualizers along the way.

jump to: why not Elasticsearch · anatomy of a line · one field, two types · sort key · bloom filters · never drop a line · query language · writing a grammar · GLOBAL IN · trade-offs

1. Why not just run Elasticsearch?

Our reason started with something boring: we already had a monitoring platform. Metrics, alerts and topology all lived on ClickHouse. Logs were the only thing with a separate stack, which meant a second cluster to run, a second query language to learn, and a second place to look during an incident. We wanted one platform, not two.

Then we looked at whether a column store could actually do the job, and at what it would cost.

Elasticsearch is excellent at what it was built for, which is full-text search. It finds text with an inverted index: for every word, a list of the documents that contain it. Looking up a word is a direct jump to that list, which is very fast.

That speed has costs, and on a log workload you pay all of them. The index can be as large as the data it indexes. The JVM heap and the shard layout need planning and re-planning. And every field has a mapping, a type decided early, often by the first document that happened to contain the field. Documents that disagree are rejected. New JSON keys keep adding fields to the cluster's state, until you hit a "mapping explosion" and start setting limits on how many fields a log line is allowed to have.

Then look at what log queries actually do. They almost always have a time range. They are usually narrowed by host or file. And very often they count and group: errors per host per minute, top paths by status, how many lines matched. Ranking by relevance almost never comes up. Nobody wants the "best" log line. They want the ones from 11:14 on api-07, in order.

That is a column store's favorite kind of query.

↳ two shapes of log store

Search engineColumn store
Finds text withInverted indexSkip indexes + scanning a narrow range
SchemaPer-field mapping, decided earlyFixed core + flexible attributes
Type conflictRejected or droppedNot a problem (section 3)
AggregationsGoodExcellent
Relevance rankingYesNo (time order)

The cost. Columnar storage compresses repetitive log data far better than an inverted index does, and it does not need a large JVM heap on every node. Teams that moved logs off Elasticsearch have published the results: Didi cut its observability hardware cost by over 30%, Uber cut the hardware cost of its log platform by more than half while serving more production traffic, and Cloudflare saw each log go from about 600 bytes as an Elasticsearch document to about 60 bytes as a ClickHouse row.

To be fair, Elasticsearch has narrowed the storage gap. Its newer logsdb index mode cuts log storage substantially. The shape of the queries still favors a column store, though, and for us, one platform beat two.

Neither column is "better". They are two shapes. The rest of this post is about what it takes to make the column store shape work for logs, where the hard part is not speed. It is that logs refuse to have a schema.

2. Anatomy of a log line

The first idea is to stop treating all fields equally. Every log line has a few things that are always there, and a long tail of things that are sometimes there. Store them differently.

  • The envelope is a handful of real columns that every line has: the time, the host, the file it came from, the severity, the original body, and how it was parsed. Keep this list short. Every column here should be one you sort by or filter on constantly.
  • The attributes are everything else, stored as key/value pairs. A service can add a field tomorrow and nothing in the table has to change.
  • The body is always kept. Parsing adds structure next to the original text. It never replaces it. If the parser is wrong, the evidence is still there.

Here is one JSON line, split that way:

simplified illustration
{"level":"error","msg":"payment failed","order_id":"A-1029","amount":499.0,"retry":true}

envelope:    severity=ERROR   body=<original line>
strings:     msg="payment failed", order_id="A-1029"
numbers:     amount=499.0
booleans:    retry=true

The same idea works for logfmt, syslog and access-log lines: each has its own small parser, and each produces the same envelope plus attributes. Nested JSON is flattened into dotted keys, so {"http":{"status":502}} becomes http.status = 502. Try it with the four common shapes, or paste a line of your own:

↳ log line splitter

body = {"level":"error","msg":"payment failed","order_id":"A-1029","amount":499.0,"retry":true,"http":{"status":502}}

envelope (real columns)

  • —

strings map

  • —

numbers map

  • —

booleans map

  • —

Hover a field to see where it came from. Edit the line and break it on purpose: a line that will not parse still lands, as parsed_as=plain text, with the body intact. The original line is always kept.

You may have noticed that attributes land in three different places, not one. That is not decoration. It is the answer to the page at the top of this post.

3. One field, two types

Back to the incident. Service A logs status: 200. Service B logs status: "OK". A mapping-based store must pick one type for the field called status, and whichever it picks, one service loses.

The fix is to stop deciding the type per field name. Decide it per value, every time a line arrives. Keep one key/value map per type: a number goes into the numbers map, a string into the strings map, true/false into the booleans map.

simplified illustration
-- one map per value type
strings   Map(String, String),
numbers   Map(String, Float64),
booleans  Map(String, UInt8)

Now both lines are stored. status = 200 lives in the numbers map, status = "OK" in the strings map. Nothing is rejected, and nobody had to agree on anything in advance. The table's structure also stops growing: a thousand new keys are a thousand new map entries, not a thousand new columns.

↳ mapping conflict arena

fixed mapping

field types, locked on first sight

none yet

rejected

  • —
0
stored
0
rejected
0
fields in mapping

typed maps

three fixed columns; keys live inside them

strings
numbers
booleans
0
stored
0
rejected
3
columns, always

Send A first, then B. The left store has already decided status is a number, so B bounces. The right store puts 200 in the numbers map and "OK" in the strings map. Then press the new-keys button a few times and watch only the left side grow.

Why not one string map for everything?

It is tempting. Everything is text anyway, so store everything as a string. Three problems:

  • Comparisons go silently wrong. duration_ms > 2000 on strings compares character by character, and "999" > "2000" is true, because 9 sorts after 2. No error. Just wrong answers.
  • Every query needs a cast, on every row it reads. That is slow, and easy to forget in exactly the query that matters.
  • Sums and averages need numbers. "p95 of duration_ms per host" should be a plain numeric aggregate, not a parse of millions of strings.

With a numbers map, a numeric comparison is just a numeric comparison. And the user never sees any of this. They write status = 200 or status = "OK", and the query layer picks the right map from the type of the value they wrote.

One gotcha worth sharing: some values look like numbers and are not. order_id=007, parsed as a number, becomes 7, and now a search for order 007 finds nothing. Identifiers with leading zeros must stay strings. So must anything else where the digits are a name rather than a quantity: zip codes, phone numbers, account numbers. The rule we ended up with is boring and works: a value with a leading zero stays a string, even when every character is a digit.

4. The sort key decides everything

In a column store, the sort order decides which queries are fast. It matters more than any index you add later, so every column has to earn its place against it.

Here is how a read works, simplified. Rows are sorted by the sort key and stored in fixed-size blocks, called granules. A tiny index, the sparse primary index, stores only the first key of each block. It is small enough to always live in memory. A query looks at it to find which blocks could contain matching rows and skips all the others without reading them.

Now picture the obvious choice for logs: sort by time. The query "logs for web-01 in the last hour" can skip everything older than an hour, which is good. But inside that hour, every host's rows are mixed together, and the index knows nothing about hosts. So the query reads every host's rows for that hour and throws most of them away.

Sort by host first, then time, and each host's rows sit together. The same query jumps to web-01's slice and reads only its last hour.

simplified illustration
ORDER BY (host, timestamp)
PARTITION BY toDate(timestamp)

↳ granule skipper · 354 rows, 45 blocks of 8

api-07Sat 21:00
auth-02Sun 03:00
web-01Sun 09:00
api-07Sun 21:00
auth-02Mon 03:00
web-01Mon 09:00
api-07Mon 21:00
auth-02Tue 03:00
web-01Tue 09:00
api-07Tue 21:00
auth-02Wed 03:00
web-01Wed 09:00
api-07Wed 21:00
auth-02Thu 03:00
web-01Thu 09:00
api-07Thu 21:00
auth-02Fri 03:00
web-01Fri 09:00
api-07Fri 16:00
auth-02Fri 17:00
web-01Fri 18:00
api-07Fri 20:00
auth-02Fri 21:00
web-01Fri 22:00
api-07Sat 00:00
auth-02Sat 01:00
web-01Sat 02:00
api-07Sat 04:00
auth-02Sat 05:00
web-01Sat 06:00
api-07Sat 08:00
auth-02Sat 09:00
web-01Sat 10:00
api-07Sat 12:00
auth-02Sat 13:00
web-01Sat 14:00
api-07Sat 14:10
auth-02Sat 14:15
web-01Sat 14:20
api-07Sat 14:30
auth-02Sat 14:35
web-01Sat 14:40
api-07Sat 14:50
auth-02Sat 14:55
web-01Sat 15:00
– / 45
blocks read
–
blocks with a real match
–
skipped

Each chip is a block, labeled with the only thing the index stores: its first key. Green = read and useful, amber = read but nothing matched, faded = skipped. Try "all hosts, last 5 minutes" under both orders. Host-first is a little worse there, because the newest rows are now spread across one tail per host. That is the trade.

The trade-off is real, and the visualizer shows it. "All hosts, last five minutes" gets a little worse with host-first sorting, because the newest rows are now spread across one tail per host instead of sitting together at the end. We took that trade, because almost every real query names a host or a file, and the ones that don't are usually small time ranges anyway.

Two more ideas, one line each:

  • Partition by day. Old days are separate parts that a query for "last hour" never opens. And deleting old logs becomes dropping a whole partition, which is instant, instead of deleting rows. A TTL set to drop whole parts does it for you.
  • Use high-precision timestamps. A busy service writes many lines in the same second. With second-level timestamps, their order is lost. Sub-second precision keeps them in the order they were written.

The rule we use: if a column is not part of how you sort, and is not filtered on constantly, it belongs in the attributes, not in its own column.

5. Finding text without a search engine

This is the question everyone asks: without an inverted index, how does a search for "connection refused" avoid reading every line?

The answer is a skip index. Each block keeps a small summary of what it contains. Before reading a block, the query asks the summary: could this block contain the word? If the answer is no, the block is skipped.

For text, the summary is a token bloom filter. Split the block's text into words, and put every word into a bloom filter: a small array of bits, where each word sets a few bits chosen by hashing it. To ask about a word, hash it the same way and check those bits. If any bit is unset, the word is definitely not in the block. If all are set, it is maybe there. A bloom filter can be wrong about "maybe". It is never wrong about "definitely not".

So a search becomes two steps:

  1. Skip every block whose filter says "definitely not".
  2. Read the "maybe" blocks, and check each line for the real phrase.

↳ bloom filter sieve · 20 blocks, k = 3 hashes

filter size (bits per block)32

tokens: connection → [25, 4, 31]refused → [3, 12, 13]

block 0definitely not
block 1false positive
heartbeat ok
payment failed for order A-1941
tls handshake error from 10.0.4.5
retrying request attempt 1
block 2definitely not
block 3definitely not
block 4definitely not
block 5definitely not
block 6definitely not
block 7match
checkout completed order A-4889
connection refused to db-03:5432
checkout completed order A-6658
slow query took 3109ms
block 8definitely not
block 9definitely not
block 10definitely not
block 11definitely not
block 12definitely not
block 13definitely not
block 14match
queue depth 1645
connection refused to db-03:5432
user 5947 logged in
heartbeat ok
block 15definitely not
block 16definitely not
block 17match
disk usage at 80 percent
timeout waiting for upstream
connection refused to db-03:5432
heartbeat ok
block 18definitely not
block 19definitely not
16
skipped
4
maybe (read)
3
real matches
1
false positives
80%
bits set

Under each block is its filter. Probed bits are green when set and red when not; one red bit is enough to skip the block. Shrink the filter and the bits fill up, so more blocks say "maybe" and turn out empty. A filter never says "not here" when the word is there.

The same trick works for things that are not prose. A bloom filter over trace IDs makes "show me all the logs for this trace" fast. A bloom filter over attribute keys makes "lines that have a user_id at all" fast.

The honest limits

  • False positives exist. A busy block has more words, so more bits are set, so it says "maybe" more often. A bigger filter helps, at the cost of more space per block.
  • Common words skip nothing. error is in nearly every block. For that search, the time range and the host filter do the real work, which is one more reason the sort key matters so much.
  • Partial words cannot use the index. The filter knows refused, not efuse. A substring search still works. It is just a scan of whatever the time range and host filter leave, which is usually small enough.

6. Never drop a line

The incident at the top of this post was a pipeline that put understanding a line ahead of keeping it. We reversed that: the pipeline's first job is to store the line. Understanding it comes second.

  • If parsing fails, keep the line as plain text, and record that something was off. Later, "show me lines that came in broken" is a filter, not a mystery.
  • Size limits are a safety net, not a filter. When a line is too big, trim it, mark it as trimmed, and keep the event. A line with a cut-off stack trace is far more useful than no line.
  • Keep experiments isolated. If a new destination shares a pipeline with a working one, and the new one fails or slows down, back-pressure can stall the working one too. Give anything experimental its own path. We wrote about how back-pressure moves through a pipeline in the Vector diff post.
  • Test with real, messy lines, not tidy JSON. A sample of real production lines caught bugs that clean test data never would have.

↳ real lines that broke clean parsers

  • "ts": 1760094902123
    before: read as seconds: a date tens of thousands of years from nowafter: detect seconds vs milliseconds by magnitude
  • order_id=007
    before: stored as the number 7after: values with a leading zero stay strings
  • level=NOTICE
    before: fell through to UNKNOWNafter: map NOTICE to INFO
  • level=crit
    before: fell through to UNKNOWNafter: map short forms: crit, err, warn
  • Warning: disk 91% full
    before: no level field at all, so UNKNOWNafter: look for a severity word at the start of the body
  • {"msg":"ok", "user":{
    before: broken JSON, line droppedafter: kept as plain text, marked as a parse failure

None of these are clever. All of them were found by feeding the parser what services actually write, instead of what we assumed they write.

7. A query language humans can read

The people searching logs at 2 AM are on-call engineers, not SQL experts, and they are in a hurry. We wanted queries that read like a sentence. One example, as a flavor of the idea rather than a spec:

simplified illustration
find "connection refused"
where severity >= error
last 1 hour
count by host

A few principles held the design together:

  • Clauses start with plain words that say what they do.
  • Logic is always written out. and and or are words. A space never secretly means "and".
  • Users write friendly field names and never see how things are stored, including which typed map a value lives in.
  • Errors point at the exact mistake and suggest a fix: "did you mean severity?"

The language is the visible part. The architecture behind it is the real lesson:

↳ where each step runs

textbrowser
→
parsebrowser
→
structured querybrowser
──▶
validateserver
→
compile to SQLserver
→
databaseread-only

Only the structured query crosses the network. User text never reaches the server.

  • Parse in the browser. Errors appear while you type, and autocomplete comes from the same grammar that does the parsing.
  • The server never parses user text. It accepts only a strict structured format, validates every part of it, and builds SQL with values passed as parameters, never pasted into the query string.
  • Every query runs read-only, with limits on time and memory. A bad query can be slow. It cannot write anything, and it cannot take the database down with it.

↳ query journey · simplified illustration

1 · textbrowser
find "connection refused"
where severity >= error
last 1 hour
count by host
2 · parse treebrowser
…
3 · structured querysent to server
…
4 · SQLserver, after validation
…

Same color, same meaning, in every panel. Stages 1 and 2 happen as you type. Only stage 3 crosses the network, and the server checks it again before building SQL. With the typo, the journey ends in the editor.

8. Writing a grammar (and why we moved parsing to the browser)

A query language is a grammar plus a parser. That sounds far scarier than it is. And as it turned out, where the parser runs mattered as much as how it is written.

Every grammar in this section is a toy teaching language. It is small on purpose, it is not our real one, and that is fine: the goal is the method.

8.1 Our first plan: do it all on the server

The first design was simple. The browser sends the raw query text. One backend compiler does everything: split the text into words, build a tree, check it, and turn it into SQL. One place, one language, one codebase.

Then we thought through what the editor actually needs, and it got expensive fast.

  • Every keystroke becomes a network call. Red underlines, syntax colors and autocomplete all need a parse. That means a request per keystroke, or an editor that lags behind your typing.
  • The server pays for half-typed queries. Most parses would be of text the user has not finished writing.
  • Two parsers, drifting apart. The editor has to tokenize the text anyway, to color it. So the browser and the server each end up with their own idea of the language, and they slowly start to disagree.
  • A hand-written parser grows messy. Every new keyword touches the tokenizer, the parser and the error messages. Good error positions, "the mistake is here", are easy to lose along the way.
  • More attack surface. The server would be running a parser on raw, untrusted text from anyone.

8.2 What we did instead

We wrote the language down as a grammar file and let a parser generator build the parser from it. We used Lezer, the parser system behind the CodeMirror editor. The generated parser runs in the browser, and one grammar drives syntax colors, autocomplete, error underlines and the structured query sent to the server.

The server shrinks to three jobs: validate the structured query, compile it to SQL, and run it safely. Lezer also parses incrementally: when you type one character, it re-parses only the part of the tree that changed, so the editor stays fast even on long queries.

↳ before: the server does everything

editorraw text, every keystroke
──▶
serverlex + parse + check + compile
→
database

↳ after: the browser understands, the server checks

editorgrammar: lex + parse + errors
──▶
servervalidate + compile
→
database

One request, on Run. The server still re-checks everything it receives.

8.3 How to write a grammar: lexer, then parser

Take one toy query and follow it through the two layers:

toy language
where severity >= "error" and host = "web-01"

Step 1: the lexer, or tokenizer, turns characters into words, called tokens. It does not care about meaning, only shapes: this run of letters is a word, this is an operator, this is a quoted string.

tokens
[KEYWORD where] [FIELD severity] [OP >=] [STRING "error"] [KEYWORD and] [FIELD host] [OP =] [STRING "web-01"]

Whitespace is thrown away here, so the parser never has to think about spaces.

Step 2: the parser checks that the tokens come in a legal order and builds a tree:

tree
Where
└── And
    ├── Compare  severity >= "error"
    └── Compare  host = "web-01"

Step 3: write the rules. Start with plain EBNF on paper. Read each = as "is made of", { } as "zero or more", and | as "or":

toy grammar · EBNF
query    = clause { clause } ;
clause   = find | where | last ;
find     = "find" STRING ;
where    = "where" expr ;
last     = "last" DURATION ;

expr     = term { "or" term } ;          (* or binds loosest  *)
term     = factor { "and" factor } ;     (* and binds tighter *)
factor   = "not" factor
         | "(" expr ")"
         | compare ;
compare  = FIELD OP value ;
value    = STRING | NUMBER ;

Step 4: the same toy grammar for a parser generator. This is Lezer-style and illustrative only, so check the syntax against the library's docs before you build on it:

toy grammar · lezer-style, illustrative
@top Query { clause+ }

clause { Find | Where | Last }

Find  { kw<"find"> String }
Where { kw<"where"> expr }
Last  { kw<"last"> Duration }

@precedence { not, and @left, or @left }

expr {
  Compare |
  Not { !not kw<"not"> expr } |
  And { expr !and kw<"and"> expr } |
  Or  { expr !or  kw<"or">  expr } |
  "(" expr ")"
}

Compare { Field Op value }
value   { String | Number }

kw<term> { @specialize[@name={term}]<Field, term> }

@skip { space }

@tokens {
  Field  { $[a-zA-Z_] $[a-zA-Z_.0-9]* }
  String { '"' (!["\\] | "\\" _)* '"' }
  Number { @digit+ }
  Duration { @digit+ $[mhd] }
  Op     { "=" | "!=" | ">=" | "<=" | ">" | "<" }
  space  { @whitespace+ }
}

Four ideas carry most of the weight for a beginner:

  1. Keywords are words with a special meaning. kw<…> takes something that would be a Field and marks it as a keyword when the text matches exactly. That way where is a keyword, but where_clause is still an ordinary field name.
  2. Precedence decides the shape of the tree. Without it, a = 1 or b = 2 and c = 3 is ambiguous. With and binding tighter than or, it always means:
    Or
    ├── a = 1
    └── And
        ├── b = 2
        └── c = 3
    It is the same as 1 + 2 × 3 in arithmetic. Write precedence down in the grammar, so the language never quietly changes its mind.
  3. Error recovery. A good parser does not give up at the first mistake. It marks the broken spot and keeps going, so the editor can still color and understand the rest of the line:
    where severity >= and host = "web-01"
                      ^ expected a value after ">="
  4. Every node knows its position. Each tree node carries its start and end character. That is what makes exact red underlines, hover hints and "did you mean" possible.

The last step: tree to structured query. A small function walks the tree and builds the plain object the server receives. The tree is a detail of the browser. The structured query is the contract between browser and server. In the toy language, it could be as small as this:

toy structured query
{ "text": "connection refused", "range": "1h" }

The workbench below runs the toy language live. It uses a small hand-written parser rather than Lezer, so you can see every moving part, but it behaves the same way: tokens, a tree with positions, errors that do not stop the parse, and a precedence switch.

↳ grammar workbench · toy language

where a = 1 or b = 2 and c = 3
precedence
parsing runs in

1 · tokens

keyword wherefield aop =number 1keyword orfield bop =number 2keyword andfield cop =number 3

2 · tree (start–end)

Query 0–30
└ Where 0–30
└ Or 6–30
└ a = 1 6–11
└ And 15–30
└ b = 2 15–20
└ c = 3 25–30

3 · errors

✓ no errors

where clause matches 5 of 8 sample rows

browser
server
requests sent: 0delay before the underline appears: 0 ms

Load the precedence preset and flip the toggle: the tree reshapes and the match count changes, from the same text. Then switch parsing to the server and type a few characters. Every keystroke is a round trip, and the highlighting lags behind your typing.

8.4 Decisions that held up

↳ what we would do again

  • Grammar as a file, not as code. Adding a keyword means changing a few lines, not touching three layers of hand-written parser.
  • One grammar, many features. Colors, autocomplete, errors and query building cannot drift apart, because they all come from the same file.
  • Explicit and. A space never secretly means "and". It keeps the grammar simple and the queries easy to read.
  • Precedence written down early, before anyone had saved queries that depended on it.
  • Positions on every node, from day one. Retrofitting good error messages later is painful.
  • A versioned structured format between browser and server, so the language can grow without breaking saved queries.
  • The server still never trusts the browser. Parsing moved to the client. Validation did not. Anything can send a request, so the server checks every field, operator and limit again.

9. The Distributed table trap

This one is general ClickHouse knowledge, and it bit us anyway: a query that works on one server can be refused, or worse, silently wrong, on a cluster.

The setup: data is split across shards. A Distributed table is a view over all of them. Below, logs is that Distributed table and logs_local is the table on each shard. A query against it is sent to every shard, and the partial results are merged on the server that received the query, the initiator.

The query: "errors over time for the top 10 hosts". The natural SQL finds the top 10 in a subquery and filters with host IN (SELECT … top 10 …). On a single server, it works.

On the cluster, ClickHouse refuses it, with error 288: double-distributed subqueries are denied. The reason is the fan-out. Each shard would run the inner query against the Distributed table, which sends it to all shards again. N shards means N × N queries. With 4 shards that is 16. With 20, it is 400, for one chart.

The tempting workaround is worse. Point the subquery at the local table instead. The error goes away. But now each shard picks its own top 10 from only its own data, and the merged chart shows the wrong hosts, or the right hosts with partial counts. There is no error at all. It just looks plausible.

The fix is GLOBAL IN. The subquery runs once, on the initiator. Its small result, ten host names, is sent to every shard as a temporary table, and every shard filters with the same list.

simplified illustration
WHERE host GLOBAL IN (
    SELECT host FROM logs
    WHERE severity = 'ERROR' AND timestamp > now() - INTERVAL 1 HOUR
    GROUP BY host ORDER BY count() DESC LIMIT 10
)

↳ shard fan-out · top 5 hosts, 4 shards

shards4

WHERE host GLOBAL IN (SELECT host FROM logs …)

subquery runs once here[5][5][5][5]initiatorshard 1shard 2shard 3shard 4
1
subquery executions
correct
what you see

real top 5

db-04794
api-09745
api-07652
edge-01609
web-03526

what the chart shows

db-04794
api-09745
api-07652
edge-01609
web-03526

IN: every shard would run the inner query against every shard again, so the count grows as shards². ClickHouse refuses it. Local IN: each shard trusts its own top list, so the merged chart can show the wrong hosts (red) or partial counts (amber), with no error. GLOBAL IN: the subquery runs once, and the small list is shipped to every shard.

The lesson: test on a real multi-shard setup, even a tiny one. Two shards in containers on a laptop would have caught both the error and the silent wrong answer. One shard catches neither.

Trade-offs and takeaways

None of this is free. What you give up compared with a search engine:

  • No relevance ranking. Results come back in time order, which for logs is usually what you want anyway.
  • Common-word searches do not skip much. The time range and host filter have to carry them.
  • No substring index. Partial-word search is a scan of the narrowed range.
  • A custom query language is real work to build, document and maintain.
  • Field discovery has to be built. A search engine gets the list of fields for free from its mappings. With typed maps, you collect and serve it yourself.

What we would tell anyone starting the same project:

  • Fixed core columns, flexible attributes for everything else.
  • Pick the type per value, not per field name.
  • Sort by what almost every query filters on.
  • Use bloom-filter skip indexes instead of an inverted index.
  • Never drop a line. Mark it and keep it.
  • Parse queries in the client, and trust only validated structured input on the server.
  • Write the language as a grammar file and let a parser generator do the heavy lifting.
  • Use GLOBAL IN for subqueries on Distributed tables, and test on more than one shard.

Logs don't have a schema. Stop asking them for one.


Related: We diffed our pipeline against Vector's source covers the pipeline that feeds a store like this, and why back-pressure decides its throughput. Kafka beyond the basics covers the log that usually sits in front of it.