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.
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
- Watch
- Try it
- Predict
- Capture
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.
Highlighted lines are the ones running in the diagram right now.
def run(q, start, end):if would_fetch_too_many_series(q, start, end):refuse(max_fetched_series_per_query)# split_queries_by_interval = 24hpieces = 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)
def get_or_miss(q, pieces):hits, misses = [], []for piece in pieces:entry = None# cache_results is off by default; freshness 10mif 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
def execute(q, piece):series = []for part in shard(q, n = 16): # on by default# query_ingesters_within = 13hif piece.end > now() - 13h:series += ingesters.select(part, piece)# query_store_after = 12h, windows overlapif 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