Split, shard, and cache the past

A query frontend splits a long-range query into day-sized pieces, spreads them across workers, merges them and caches finished days, which almost never change, so a refresh recomputes only the newest slice.

Previously

Years of data now sit in compacted day-sized blocks; a dashboard still has to read 30 days of them quickly.

Scene 12

Split, shard, and cache the past

  1. Watch
  2. Try it
  3. Predict
  4. Capture
instrumentcollectstorequeryone jobalertnotifyDASHBOARDlast 30 dayssum(rate(http_requests_total[5m])) by (route)sum(rate(http_requests_total[5m])) by (route)QUERY FRONTENDno splitting · no spreading · no cacheidle: one query in, the same one queryout, a single job over the whole rangePIECES · 1 SLICEoldest ← → nowp1: both, idlenewer than 13 h → ingesters · older than 12 h → blocks · the 1 h overlap is read from bothnewer than 13 h → ingesters · older than 12 h → blocks · the 1 h overlap is read from bothfrom blocksfrom ingestersoverlap: bothcachedrecomputescanningWORKERS0 of 1 busyw1RESULTS CACHEOFFopt-in: off unless you enable itopt-in: off unless you enable itcache off: every refresh recomputes every pieceopt-in — not enabled hereLATENCY—series touched: ~2,500series touched: ~2,500A 30-day dashboard query arrives. Nothing has cut it up yet: one job is about to scan every day in the range.
What to watch for

How does a dashboard asking for 30 days of data come back fast? Right now it doesn't. Watch one panel's query — requests per second per route, over the last 30 days — arrive as a single job. One worker walks the whole range: the last few hours still live in the ingesters' memory, and everything older is read from the compacted day blocks you just built. The bar at the bottom is the wall-clock for the whole thing. Notice the box the query passes through on its way in: it is doing nothing yet, and it is the thing this scene turns on. Its name is the query frontend — the service that sits in front of the workers and is allowed to rewrite a query before anyone runs it.

Continue unlocks when the animation finishes.
Implementation

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

QueryFrontend.run
rewrites one long query into pieces and merges the answers
def run(q, start, end):
if would_fetch_too_many_series(q, start, end):
refuse(max_fetched_series_per_query)
# split_queries_by_interval = 24h
pieces = split(align(start, end, q.step), 24h)
hits, misses = cache.get_or_miss(q, pieces)
fresh = workers.run_parallel(q, misses)
cache.put(q, fresh)
return merge(hits + fresh)
Cache.get_or_miss
decides which pieces the shelf is allowed to serve
def get_or_miss(q, pieces):
hits, misses = [], []
for piece in pieces:
entry = None
# cache_results is off by default; freshness 10m
if cache_results and piece.end <= now() - 10m:
entry = memcached.get(key(q, piece, q.step))
if entry:
hits.append(entry.answer)
else:
misses.append(piece)
return hits, misses
Querier.execute
one worker runs one piece against ingesters and blocks
def execute(q, piece):
series = []
for part in shard(q, n = 16): # on by default
# query_ingesters_within = 13h
if piece.end > now() - 13h:
series += ingesters.select(part, piece)
# query_store_after = 12h, windows overlap
if piece.start < now() - 12h:
series += blocks.select(part, piece)
if chunks(series) > max_fetched_chunks_per_query:
abort("too many chunks")
return eval(q, series)

Where this sits in Metrics / Monitoring System

Scene 12 of 18, in the Scale out act — Remote-write, ingesters, dedup, blocks, split queries.. A query frontend splits a long range into day-sized pieces, spreads them across workers and caches finished days, which almost never change, so a refresh recomputes only the newest slice.

Up next. Now that query cost follows the number of series a query touches, the next question is why that number can explode even while the dashboard's count of live series stays flat.

All 18 scenes in Metrics / Monitoring System · Every curriculum

Built with Arqly
Every scene in Metrics / Monitoring System builds on the one before it.All 18 Metrics / Monitoring System scenes