Same query, two execution plans

ELK answers service=api ERROR last 1h by intersecting posting lists in the inverted index — milliseconds; Loki answers it by resolving label matchers to chunk refs, fetching chunks from object storage, and grepping them in the querier process — seconds to minutes. Same result, opposite ends of the same trade-off curve.

Previously

Segments-with-inverted-index on one side, chunks-plus-tiny-index on the other. Same input, different disk shape — now run the same question through both and watch the execution plans diverge.

Scene 05a

Same query, two execution plans

  1. Watch
  2. Try it
  3. Predict
  4. Capture
QUERYservice=api AND level=E…window: 1h(no content filter)ELKelasticsearch / kibana1parse KQL{ }2resolve time-range index…[ ]3scatter to shards · post…service…121927344153677889level=E…9192741586781↓ sorted-merge intersectresult192741674fetch top-N source docstop-N _sourcems2 ms3 ms28 ms7 msTOTAL40 msLOKIgrafana loki1parse LogQL{ }2resolve label matchers →…{service=api, level=error}chunk_refs:603fetch chunks from S3S3slow× 604decompress + line-filter…decompress block-by-blockms4 ms9 ms2.2 s4.1 sTOTAL6.3 sBaseline 1-hour window. ELK finishes in tens of milliseconds — stage 3 (posting-list intersection) dominates.…
What to watch for

One query: service=api AND level=ERROR over the last hour. Both systems return the SAME result set — the difference is the work each one does to get there. We call that ordered set of stages an execution plan (or read path). Watch ELK's plan run across the top, then Loki's plan across the bottom. Note where the time goes in each lane.

Continue unlocks when the animation finishes.
Implementation

Highlighted lines are the ones running in the diagram right now.

ELK.search
coordinating node fans out, posting lists do the filtering
def search(kql, timeWindow):
ast = parseKQL(kql)
indices = resolveIndexPattern(
alias='logs-*', window=timeWindow,
)
shards = scatter(ast, indices)
perShard = []
for shard in shards:
# sorted-merge intersect across query terms
docIds = intersect_postings(shard, ast.terms)
topN = score_bm25(shard, docIds)[:N]
perShard.append(topN)
merged = mergeTopN(perShard)
return fetchSourceDocs(merged)
Loki.query
labels resolve to chunk_refs; the regex runs in the querier
def query(logql, timeWindow):
ast = parseLogQL(logql)
chunkRefs = index.lookup(
labelMatchers=ast.matchers,
window=timeWindow,
)
results = []
for ref in chunkRefs:
chunk = objectStore.get(ref) # S3
for block in chunk.decompress():
for line in block:
if ast.lineFilter.match(line):
results.append(line)
return results
Lucene.intersect_postings
galloping sorted-merge, smallest list drives
def intersect_postings(shard, terms):
lists = [shard.postings(t) for t in terms]
lists.sort(key=len) # smallest drives
out = []
for candidate in lists[0]:
keep = True
for other in lists[1:]:
# gallop forward to >= candidate
if not other.advanceTo(candidate):
keep = False
break
if keep:
out.append(candidate)
return out

Where this sits in Build a distributed logging stack (ELK / Loki)

Scene 05a of 12. ELK answers via posting-list intersection — milliseconds; Loki resolves labels to chunks, fetches them from S3, and greps in-process — seconds to minutes. Opposite ends of the same trade-off curve.

Up next. ELK answered in milliseconds because the index already knew which documents contained ERROR. The temptation is obvious — just index more things. The next scene shows what happens when 'more things' has unbounded values.

All 12 scenes in Build a distributed logging stack (ELK / Loki) · Every curriculum

Built with Arqly
Every scene in Build a distributed logging stack (ELK / Loki) builds on the one before it.All 12 Build a distributed logging stack (ELK / Loki) scenes