Time-Series Storage & Roll-ups
The Time-Series Database That Froze on Flush
Test your architecture intuition: Pitch a 7-axis solution, survive two aggressive reviewer objections, and inspect the staff-level Teacher Gold Answer.
Part 0. Start here
The problem: two questions about one metric
Basketly, the grocery shop of pages 10, 12 and 14, runs its checkout service on four pods. Three of them take 20 requests a second each. The fourth, p4, is a canary: it runs the next build and takes 2 requests a second, 1 request in 31. Every pod's latency is sampled every 10 seconds, and the store rolls the samples up to 1-minute and 1-hour resolution.
On Tuesday at 10:20:00 the canary gets a slow build. For 90 seconds, every request it serves takes 2 seconds, until the rollback at 10:21:30. Two questions follow.
Question 1, the incident review on Wednesday: "What were checkout's mean latency and p99 during those 90 seconds?" The team's dashboard averages each pod's mean and each pod's p99.
| What was asked | The dashboard (averages of the four pods) | The truth (all requests together) | How wrong |
|---|---|---|---|
| Mean latency | (3 × 0.065 + 2.0) ÷ 4 = 0.549 s | (60 × 0.065 + 2 × 2.0) ÷ 62 requests a second = 0.127 s | 4.3× too high |
| p99 | (3 × 0.4 + 2.0) ÷ 4 = 0.8 s | 2.0 s: the canary's slow requests are 2 in 62, 3.2% of all requests, more than 1% | 2.5× too low |
Averaging broke both numbers, in opposite directions. The mean went up because the canary's 2 requests a second counted as much as another pod's 20. The p99 went down because three healthy pods outvoted the one that was slow.
Question 2, the capacity review a month later: "Show the yearly p99, with that incident on it." A yearly panel reads the 1-hour tier (Part 8), so what it shows for that morning is the 1-hour roll-up of 10:00 to 11:00.
| What the hour can say | Arithmetic | Result |
|---|---|---|
| The incident's share of the hour | 180 slow requests of 62 × 3,600 = 223,200 | 0.08% |
| The hour's true p99 | 0.08% is far under 1% | 0.4 s, as in any normal hour |
| An hour built by averaging each minute's p99 | 58 minutes at 0.459, 10:20 at 2.035, 10:21 at 1.570 | 0.503 s: neither the hour's p99 nor the incident's |
| The p99 estimated from the hour's summed bucket counts | Part 6 | 0.462 s |
| What still shows the incident | the (1, 2.5] bucket's count, and the in-flight gauge's maximum | 180 requests, and 4 in flight on p4 (normally 0) |
Two facts frame the page. The Gorilla paper, whose compression Prometheus and many other time-series stores use, starts from "16 bytes per point" (an 8-byte timestamp and an 8-byte value) and reports that it can "compress time series to an average of 1.37 bytes per point, a 12x reduction in size". And Amazon CloudWatch keeps sub-minute points for 3 hours, 1-minute points for 15 days, 5-minute points for 63 days and 1-hour points for 455 days: "Data points that are initially published with a shorter period are aggregated together for long-term storage". Every long-lived metrics store rolls up. The question is what it keeps when it does.
You must keep checkout's latency for five years, graph any range in a second, and answer "what was the p99?" for any minute, hour or month. You can't keep every sample forever. What do you keep for each old minute and hour, so the average and the p99 of any range can still be computed? And what can't you keep, however much you are willing to pay?
The big picture
Synthesizing vector architecture diagram...
What to notice: samples enter at the head and are kept raw, in the head and then in blocks, for 15 days. The 1-minute tier is built from raw samples and the 1-hour tier from 60 one-minute rows, so both keep only values that merge. Queries pick a tier by their step (the three dotted arrows: raw, 1-minute or 1-hour); alert rules read raw data, because coarse tiers see an incident late or not at all.
What you'll be able to do after this page
- Count the series a metric creates, and tell a counter, a gauge and a histogram apart (Part 1).
- Explain why series, not samples, size memory and cost, and where to stop a label that multiplies them (Part 2).
- Trace a sample from the write-ahead log into a chunk, a block, and out by retention (Part 3).
- Compress timestamps with delta of deltas and values with XOR, by hand, and say when compression gets worse (Part 4).
- Choose each roll-up's aggregates by metric type, handle counter resets, and write roll-ups that are safe to run twice (Part 5).
- Compute a percentile over any range from merged buckets or a sketch, and state its error (Part 6).
- Size each tier by the values it keeps, and see through the "60× smaller" claim (Part 7).
- Pick the tier for a query, align its steps, take rates across restarts, and handle stale series and gaps (Part 8).
- Decide what happens to a sample that arrives after its minute was rolled up (Part 9).
- Work out when an alert can fire on raw data and on each tier (Part 10).
- Keep the roll-up pipeline correct when a job is retried, a worker falls behind or an ingester crashes (Part 11).
- Follow one sample through every layer, from the client library to a year-long graph (Part 12).
- Map all of it to AWS, and name the look-alikes (Part 13).
You may have arrived from a step that relies on this: steps 1.1, 1.2 and 2.1 to 2.5 of the metrics and alerting loop (chunks, the label index, Gorilla compression, roll-up tiers, cardinality, late samples) and its step 3.4 (tiers and the results cache); R1.8 of the ad-click aggregation loop, which sends readers to the metrics loop "for time-series storage in general (compression, rollups, retention tiers)"; step 3.5 of the URL shortener (unique counts that merge); or step 3.2 of the gaming leaderboard (a histogram whose error is its bucket). In all, 9 of the 38 interview loops depend on a time-series or roll-up mechanism in at least one step.
Part 1. A series is a name and labels
Four pods, one latency. How many things does the store keep? It isn't one ("checkout latency"), and it isn't four (one per pod). This Part names what the store actually keeps, and sets up the example.
Name plus labels
A series is a metric name plus one full set of label values. checkout_inflight{pod="p1"} and checkout_inflight{pod="p2"} are two series; change any label value and you have a different series. A sample is one (timestamp, value) pair of one series: a timestamp in milliseconds and a 64-bit floating-point value.
Counter, gauge, histogram
| Type | What it holds | How it changes | How you read it |
|---|---|---|---|
| Counter | A running total: requests served, seconds spent | Only goes up, except that it restarts from 0 when its process restarts | As a rate or an increase over a window, never as its raw value |
| Gauge | A reading at one instant: requests in flight, memory used | Up or down | As its value at an instant, or as min, max and average over a window |
| Histogram | One counter per bucket bound le ("less than or equal"), each counting the requests at or below that bound (cumulative), plus _sum and _count; the last bound +Inf always equals _count | Every counter only goes up | Sum the buckets' increases over a range, then estimate a percentile (Part 6); mean = increase of _sum ÷ increase of _count |
A fourth type, the summary, computes quantiles inside the client and exports them finished. Part 6 shows why two summaries can't be combined.
Here is what p4 exposes when it is scraped at 10:00:00, in the text format agents read:
textcheckout_duration_seconds_bucket{pod="p4",le="0.1"} 7200 checkout_duration_seconds_bucket{pod="p4",le="0.5"} 7200 checkout_duration_seconds_bucket{pod="p4",le="1"} 7200 checkout_duration_seconds_bucket{pod="p4",le="2.5"} 7200 checkout_duration_seconds_bucket{pod="p4",le="+Inf"} 7200 checkout_duration_seconds_sum{pod="p4"} 287.9999999999997 checkout_duration_seconds_count{pod="p4"} 7200 checkout_inflight{pod="p4"} 0
Eight lines, eight series: seven counters from one histogram, and one gauge. (_sum should be 7,200 × 0.04 = 288; the client adds 0.04 in 64-bit floating point, 7,200 times, and ends a hair below it. Part 4 shows what that costs.)
The example we follow
A real metrics platform holds a hundred million series, too many to watch. So one story runs through the whole page: one metric at 10 s, 1 minute and 1 hour resolution. It is Basketly again (pages 10, 12 and 14), now through its checkout service. Every trace on this page comes from running a private reference simulation of this setup (the requests, the client counters, the scrapes, the store, the encoders, the roll-up jobs and the queries, with every rule as a setting), not from working it out by hand. The example's numbers are small so that every value fits on a line; the last column says what changes at real scale.
| Setting | Our example | At real scale |
|---|---|---|
| Pods and traffic | p1, p2, p3: 20 requests a second each; p4 (canary): 2 a second. 62 a second, 223,200 an hour while p4 runs | Metrics loop Round 2: 100 million active series, 10 million samples a second |
| Request timing | p1 to p3: request j of each second starts at 25 ms + 50 ms × j (j = 0 to 19); p4: at 250 ms and 750 ms. No randomness | |
| Latency | p1 to p3: j = 0 to 17 take 0.04 s, j = 18 takes 0.18 s, j = 19 takes 0.4 s, so each second's mean is (18 × 0.04 + 0.18 + 0.4) ÷ 20 = 0.065 s. p4: 0.04 s | |
| The incident | Every p4 request started from 10:20:00 to 10:21:29.999 takes 2.0 s: 2 × 90 = 180 requests | |
| Metric 1: histogram | checkout_duration_seconds with bounds 0.1, 0.5, 1, 2.5 and +Inf, plus _sum and _count: 7 series a pod, 28 | Metrics loop: about 1,000 series per service |
| Metric 2: gauge | checkout_inflight{pod}, requests in flight at the scrape instant: 4 series | |
| Series | 32 | 100 million (Round 2), 1 billion (Round 3) |
| Counters start | Every pod's process starts at 09:00:00, so at 10:00:00 p1's _count is 20 × 3,600 = 72,000 and p4's 7,200 | |
| Scrapes | Every 10 s, exactly at :00, :10, …; a sample reaches the store 1 s after its timestamp | Metrics loop: 10 s. Prometheus's default scrape_interval is 1 minute |
| Head and blocks | Chunks of 120 samples; 2-hour blocks; a write-ahead log | Prometheus: "two-hour blocks", 120 samples a chunk, WAL segments of 128 MB |
| Out-of-order window | 10 minutes (variants: 0, Prometheus's default, and 60 s, Managed Prometheus's default) | Metrics loop: 10 minutes |
| Roll-ups | Minute M is rolled up at M + 1 min 30 s, from raw samples; hour H at H + 1 h 5 min, from its 60 minute rows | Thanos: 5-minute roll-ups after 40 hours, 1-hour after 10 days |
| Retention | Raw 15 days; 1-minute 90 days; 1-hour 2 years | Metrics loop: raw 30 days, 5-minute 180 days, 1-hour 5 years |
One simplification is labelled: our client library counts a request when it starts (a real one counts it when it completes), so every scrape's increase is an exact multiple. Side row G in Part 14 shows what counting at completion changes.
We follow Tuesday from 10:00:00 to 12:00:00, one 2-hour block, and then let the tiers age over days and months with arithmetic. Four things happen on that main line, spaced so they never overlap: the canary incident (10:20:00 to 10:21:30), a restart of pod p2 (10:40:03), a network partition that makes p3's samples seven minutes late (11:05:00 to 11:12:05), and the canary's removal (11:30:00). Everything else (a user_id label, a bucket layout that ends too low, samples resent newest first, a retried roll-up job, a stalled roll-up worker, a crashed ingester) runs on a copy of the main line, one per side row, so no event swamps another.
The story in twelve beats: 32 series from two metrics (Part 1); the label that made 280,000 (Part 2); a sample's trip to disk (Part 3); 16 bytes into one or two (Part 4); the minute the canary broke (Part 5); the p99 that averaging hid (Part 6); fifteen days, ninety days, two years (Part 7); the graph, the restart and the pod that left (Part 8); seven minutes late (Part 9); the alert the hourly tier never sent (Part 10); the job that ran twice (Part 11); one sample, end to end (Part 12). Each Part shows only its own events; the full table is in Part 14.
Events 1 and 2
| # | Time | Event |
|---|---|---|
| 1 | 10:00:00 | The first scrape of the block: 32 samples, one per series (4 pods × 7 + 4). p1's _count is 72,000; its le="0.1" bucket is 18 × 3,600 = 64,800; its le="0.5", le="1", le="2.5" and le="+Inf" are all 72,000, because every request took 0.4 s or less. p1's _sum is 3,600 × 1.3 = 4,680 s |
| 2 | 10:00:10 | Each counter grows by a fixed step: p1's _count +200, le="0.1" +180, _sum +13.0; p4's _count +20, _sum +0.8. The gauge reads 2 on p1 to p3 and 0 on p4 at every normal scrape, while the average number in flight on p1 is 20 × 0.065 = 1.3. A gauge is a reading at an instant, not a summary of the 10 seconds |
Why 2 and not 1.3? At each scrape instant, p1's two slow requests from the second before (0.18 s and 0.4 s, started 75 ms and 25 ms before the instant) are still running; its fast ones are not. A sample 10 seconds later sees exactly the same picture. The gauge can't tell you what happened between those two instants, however well you store it.
Snapshot T1, 10:00:10: p4's eight series in the head
| ID | Series | 10:00:00 | 10:00:10 |
|---|---|---|---|
| 25 | checkout_duration_seconds_bucket{pod="p4",le="0.1"} | 7,200 | 7,220 |
| 26 | …_bucket{pod="p4",le="0.5"} | 7,200 | 7,220 |
| 27 | …_bucket{pod="p4",le="1"} | 7,200 | 7,220 |
| 28 | …_bucket{pod="p4",le="2.5"} | 7,200 | 7,220 |
| 29 | …_bucket{pod="p4",le="+Inf"} | 7,200 | 7,220 |
| 30 | checkout_duration_seconds_sum{pod="p4"} | 287.9999999999997 | 288.8000000000001 |
| 31 | checkout_duration_seconds_count{pod="p4"} | 7,200 | 7,220 |
| 32 | checkout_inflight{pod="p4"} | 0 | 0 |
Each series holds its two samples in its own open chunk, which started at 10:00:00 and will close after 120 samples. IDs are given in the order the store first saw each series.
Checkout's latency on four pods, as a histogram with five buckets. How many series is that, and which of them change on every scrape?
What to remember from Part 1
- A series is one metric name plus one full set of label values.
- Counters only go up; a gauge is a reading at one instant; a histogram is counters, one per bucket.
- Count series before you count samples.
Part 2. Series, not samples, cost memory
At 10:50:00 on a copy of the main line (side row C), a developer adds a user_id label to checkout's histogram to debug one customer's slow orders. Checkout still serves 62 requests a second. Nothing about the traffic changes. Here is what does.
| Before (our 32 series) | After one hour with user_id | |
|---|---|---|
| Series | 32 | 40,000 users (each served by one pod) × 7 histogram series = 280,000 |
| Head memory at 8.3 KB a series | 32 × 8.3 KB = 266 KB | 280,000 × 8.3 KB = 2.3 GB |
| Samples a second | 3.2 | about 28,000, if the client library keeps exporting every user's series at each scrape (they do, until the process restarts) |
| Requests a second | 62 | 62 |
| As CloudWatch custom metrics | 4 metrics (the histograms only: one per pod, published as Values and Counts arrays): $1.20 a month | user_id × pod = 40,000 metrics: 10,000 × $0.30 + 30,000 × $0.10 = $6,000 a month |
8.3 KB is a planning figure, from Grafana Mimir's capacity guide: "Memory: 2.5GB for every 300,000 series in memory". The $0.30 and $0.10 are CloudWatch's prices per metric a month for the first 10,000 and the next 240,000 metrics (US East). The same traffic now costs almost 9,000 times the memory and 5,000 times the price.
Finding series: the label index
A query names series by their labels: checkout_duration_seconds_bucket{pod="p4", le="2.5"}. The store doesn't scan all its series to find them. It keeps an inverted index: for every label=value pair, a sorted list of the IDs of the series that carry it (a posting list).
| Label pair | Posting list (series IDs) |
|---|---|
__name__="checkout_duration_seconds_bucket" | 1–5, 9–13, 17–21, 25–29 (20 series) |
pod="p4" | 25, 26, 27, 28, 29, 30, 31, 32 |
le="2.5" | 4, 12, 20, 28 |
Event 3. The query intersects the three lists: only 28 is in all of them, so the store reads one series' chunks and nothing else. A regular expression (pod=~"p1|p2") unions the lists of the matching values first. The index costs memory per series, not per sample: about 200 B a series in the metrics loop's estimate, so 6.4 KB for our 32. The same idea, posting lists intersected, is how a search engine finds documents (Trie & Inverted Index).
What a series costs
A new sample for an existing series adds a byte or two to an open chunk (Part 4). A new series needs an ID, its label strings, entries in several posting lists, an open chunk and bookkeeping in memory: kilobytes, whatever its sample rate. That is why a series sampled once a minute costs about as much head memory as one sampled every second. "Samples tell you the bandwidth; series tell you the memory" (step 2.4 of the metrics loop).
A series stays in the head until the head is next cut into a block after its last sample (Part 3), so every user in side row C stays in memory for hours after their last request, not seconds. Managed Service for Prometheus counts the same way for its quota: "A series is active if a sample has been reported in the past 2 hours".
Where to stop a label
The metrics loop's step 2.4 stops cardinality in three layers; here they are in one paragraph each.
- At the gateway, by rule. Drop labels known to be unbounded (
user_id,request_id, raw URLs), and cap the number and length of labels on a series. The write is refused with400, and the team sees why. - Where the series live, by count. Each tenant has an active-series limit, split into a share per ingester (the servers that hold the head). A sample for an existing series is still accepted; a sample that would create a new series over the limit is dropped and counted. The limit sits where the memory is.
- On growth. Alert when a tenant's series double. To count distinct series cheaply, use a HyperLogLog sketch, which estimates "how many different?", not a Bloom filter, which answers "have I seen this one?" (Part 6).
An ingester's memory is shared by every tenant whose series hash to it, so one team's user_id label hurts every other team on those ingesters. That is why the limit is per tenant, not just per ingester. A limit on samples a second is a different limit, for bandwidth: it is usually a token bucket per tenant (Rate Limiting Algorithms). Which ingester holds which series is Sharding, Hot Keys & Rebalancing.
A developer adds user_id to checkout's histogram to debug one customer. Traffic doesn't change. What does?
What to remember from Part 2
- Every new label value is a new series, and each costs kilobytes of memory.
- Samples tell you bandwidth; series tell you memory and the bill.
- Limit series per tenant where the series live, and alert on their growth.
Part 3. Writing a sample: head, log, chunks, blocks
At 10:00:01 the 32 samples of the 10:00:00 scrape land at the store (event 4). Checkout sends 3.2 samples a second, 720 per series in every two hours, and the store must accept each in microseconds, never lose one it acknowledged, and later read any series' last two hours in one sweep. This Part follows a sample from its arrival until it is deleted.
Log first, then the chunk
The store's in-memory part is called the head. A sample goes through it in two steps:
- Append it to the write-ahead log (WAL), a file written only at its end, so the sample survives a crash. Prometheus writes its WAL "in 128MB segments". What "survives" means, and what fsync buys, is Write-Ahead Log, fsync & Group Commit.
- Append it to its series' open chunk in memory, compressed as it goes (Part 4). The sample is queryable at once.
A chunk holds one series' samples in time order. It is closed at 120 samples (20 minutes at our interval), and a new one is opened. Short chunks matter because a compressed chunk can only be read from its start.
Synthesizing vector architecture diagram...
What to notice: every arrow points forward. A sample is appended to the log and to a chunk, the head is cut into an immutable 2-hour block, blocks move to object storage and are merged, and retention deletes whole blocks. No arrow ever returns to edit a block in place.
Blocks, written once
Every so often the store takes the oldest two hours of the head and writes them as a block: every series' chunks for those two hours plus the label index (Part 2), in files that are never edited again. The WAL segments older than the block are then dropped: the block holds that data now.
When is the head cut? Prometheus waits until the head spans more than three hours (MaxTime − MinTime > chunkRange/2*3 in its source), then writes the oldest two-hour window. The extra hour is on purpose: the newest hour stays open for samples that arrive a little late (Part 9).
Events 4 to 6
| # | Time | Event |
|---|---|---|
| 4 | 10:00:01 | The 32 samples of 10:00:00 land: appended to the WAL, then to each series' open chunk; queryable at once |
| 5 | 10:00:00 → 12:00:00 | Each series gets 720 samples, cut into 720 ÷ 120 = 6 chunks, except p4's 8 series, which end at 11:29:50 with 540 samples in 5 chunks (four of 120 and one of 60), plus a staleness marker at 11:30:00 (Part 8). The block will hold 24 × 720 + 8 × 540 = 21,600 samples |
| 6 | 12:00:11 and 13:00:11 | At 12:00:11 the newest sample is 12:00:10 and the head, which started at 09:00:00, spans more than three hours: the block for 08:00 to 10:00 is written (our 09:00 to 09:59:50 samples, 11,520 of them), and the head now starts at 10:00. At 13:00:11, the arrival of the 13:00:10 samples, it spans more than three hours again: the block for 10:00 to 12:00 is written, and the WAL before 12:00 is dropped |
How old data leaves
Retention is simple because blocks are: when a block is older than the retention period, the whole block is deleted. Nothing is rewritten. Prometheus notes that "Expired block cleanup happens in the background. It may take up to two hours", so a dashboard can still show data a little past its retention.
In the background, compaction merges small adjacent blocks into larger ones (Prometheus: "up to 10% of the retention time, or 31 days, whichever is smaller"), which cuts the number of files a long query opens. It merges; it never edits a sample. This is a time-window policy: data sorted by time and expiring by age leaves "as whole files", which is the row for "Time series and caches that expire" in Part 5 of LSM-Trees & Compaction. That page also owns what happens when compaction falls behind and writes stall; its drill, The Time-Series Database That Froze on Flush, is answered there.
Why not a table of rows
The first design in step 1.0 of the metrics loop stores one row per sample in a relational table. Each row repeats the metric's name and labels and has index entries of its own: "roughly 50 to 100 bytes on disk", against 16 bytes of timestamp and value, and about 2 bytes after Part 4's compression. And a B-tree index keyed by time scatters one series' samples across pages, so a two-hour graph of one series reads hundreds of pages instead of six chunks (step 1.1 there).
Checkout writes 3.2 samples a second, 720 per series per block. Why append them to per-series chunks instead of inserting rows?
What to remember from Part 3
- New samples go to a log, then to an open chunk per series in memory.
- The head becomes immutable 2-hour blocks; retention deletes whole blocks.
- Compaction merges blocks; it never edits a sample.
Part 4. About a byte a sample
A sample is an 8-byte timestamp and an 8-byte value: 16 bytes. The metrics loop's Round 2 writes 10 million of them a second, 13.8 TB a day. But consecutive samples of one series are nearly identical: the timestamps are 10,000 ms apart, and a counter grows by about the same step each time. Gorilla, the compression in Facebook's paper and in Prometheus, stores only what changed. Here it is on p4's _count, which reads 7,200, 7,220 and 7,240 at 10:00:00, 10:00:10 and 10:00:20.
Timestamps: delta of deltas
Store the first timestamp in full (the paper: as a 14-bit offset from the block's start). After that, don't store the time, or even the gap since the previous sample: store how the gap changed, the delta of deltas D = (tₙ − tₙ₋₁) − (tₙ₋₁ − tₙ₋₂). With scrapes every 10 s, the gap is always 10 s, so D is almost always 0, and 0 costs one bit.
| D (the paper, in seconds) | Written as | Bits |
|---|---|---|
| 0 | 0 | 1 |
| −63 to 64 | 10 + 7 bits | 9 |
| −255 to 256 | 110 + 9 bits | 12 |
| −2,047 to 2,048 | 1110 + 12 bits | 16 |
| anything else | 1111 + 32 bits | 36 |
The paper found that "about 96% of all time stamps can be compressed to a single bit". Prometheus stores milliseconds, so it uses wider ranges (10 + 14 bits, 110 + 17, 1110 + 20, 1111 + 64): "Gorilla has a max resolution of seconds, Prometheus milliseconds. Thus we use higher value range steps".
Values: XOR with the previous value
Two nearby floating-point numbers share their sign, their exponent and their leading bits. XOR the new value with the previous one: the result is all zeros except in a few "meaningful" bits in the middle.
| XOR result | Written as | Cost |
|---|---|---|
| All zeros (the value didn't change) | 0 | 1 bit |
| The meaningful bits fit inside the previous value's window of leading and trailing zeros | 10 + the bits inside that window | 2 + window width |
| Otherwise | 11 + 5 bits for the number of leading zeros + 6 bits for the number of meaningful bits + those bits | 13 + meaningful bits |
Event 7: three samples of p4's _count, bit by bit (the paper's encoder, in seconds; the block starts at 10:00:00)
| Sample | Timestamp | Value | What is written | Bits |
|---|---|---|---|---|
| 1 | 10:00:00 | 7,200 | 14-bit offset from the block start (0), then the 64-bit value | 78 |
| 2 | 10:00:10 | 7,220 | The previous gap is the first sample's offset, 0, so D = 10: 10 + 7 bits = 9. 7,200 XOR 7,220 has 19 leading zeros, 3 meaningful bits and 42 trailing zeros, and there is no previous window: 11 + 5 + 6 + 3 = 16 | 25 |
| 3 | 10:00:20 | 7,240 | D = 10 − 10 = 0: 0, 1 bit. 7,220 XOR 7,240 has 17 leading zeros and 5 meaningful bits, which don't fit the previous window (19 leading zeros): 11 + 5 + 6 + 5 = 18 | 19 |
From the third sample on, the timestamp costs 1 bit, and the value costs 10 to 20 bits: this counter grows by 20 each time, so its XOR is never zero. (In milliseconds with the paper's ranges, sample 2's D would be 10,000 and cost 1111 + 32 = 36 bits; Prometheus stores the first gap as a variable-length integer instead, 2 bytes.)
What the paper measured, and what our series cost
On Facebook's workload, the paper reports: "Roughly 51% of all values are compressed to a single bit" (unchanged); about 30% take the 10 form, "with an average compressed size of 26.6 bits"; the remaining 19% take 11, "with an average size of 36.9 bits". In all, "1.37 bytes per data point". Prometheus documents "an average of only 1-2 bytes per sample". Both are averages over a workload. Ours is different.
Event 8: bytes per sample over the 10:00 to 12:00 block
| Series group | Paper's encoder, seconds (block header excluded) | Prometheus's XOR chunks, milliseconds (chunk headers included) |
|---|---|---|
Integer counters (buckets, _count), 24 series | 1.61 to 2.57, average 2.00 | 2.02 |
float64 _sum, 4 series | 6.99 | 6.60 |
| Gauges, 4 series | 0.27 | 0.40 |
| All 32 series | 2.41 | 2.39 |
| For comparison | 16 raw; 1.37 on the paper's workload | "1-2 bytes" in Prometheus's documentation |
Synthesizing vector architecture diagram...
What to notice: our gauges cost a quarter of a byte, because their value almost never changes. Our integer counters cost 2 bytes and the float _sum series cost 7, so all 32 series together average 2.41 bytes: worse than the paper's 1.37 (the second bar), and far better than 16 raw (the last bar).
This workload does worse than the paper for two reasons, and both are common:
- Counters that change on every sample never take the 1-bit
0path. Ofp1's 719 value XORs for_count, 704 take the10path and 15 the11path; none is zero. - A float that accumulates fractions is the worst case.
p1's_sumshould read 4,680 at 10:00:00, but after 72,000 additions of 0.04, 0.18 and 0.4 in 64-bit floating point it holds 4679.999999997908. Each new value differs from the last in dozens of low bits, so it costs about 7 bytes.
The gauges show the other end: p1's gauge reads 2 at all 720 scrapes, so 719 of its value XORs are zero, 1 bit each, plus 1 bit of timestamp.
Aligned scrapes keep timestamps at one bit
Side row J (a copy): the agent's timestamps drift by −1, 0 and +1 ms in turn. The gaps become 10,001, 10,001 and 9,998 ms, so D runs 0, −3, +3. In Prometheus's encoding, two timestamps in every three now cost 10 + 14 = 16 bits instead of 1: about 11 bits a timestamp on average (10.9 in the run over the delta-of-delta timestamps; each chunk's first two timestamps are stored whole). p1's _count grows from 2.03 to 3.25 bytes a sample. The paper's encoder in seconds would not see the jitter if timestamps are rounded to the nearest second, but it would if they are truncated: 10:00:09.999 becomes second 9, and p1's _count grows from 1.97 to 2.97 bytes a sample. Prometheus's scraper, for one, moves a scrape's timestamp by up to 2 ms onto its intended schedule, "to enable better compression at the TSDB level" (a comment in its source).
Samples arrive every 10,000 ms and a counter grows by 20 each time. What is the smallest thing you could write for the next sample?
What to remember from Part 4
- Store the change of the gap, and the XOR with the previous value: most samples cost a few bits.
- 1.37 bytes was one company's workload; counters and float sums cost more.
- Aligned scrapes keep timestamps at one bit.
Part 5. Roll-ups keep what can be merged
At 10:21:30 the 1-minute roll-up of minute 10:20 runs: the minute in which the canary broke. From now on, a query for "last year" should never have to decode this minute's raw samples; in 15 days they are deleted anyway. So the roll-up must store, for this one minute, everything a query a year from now could ask about it: its mean latency, its slowest moment, its p99, its request rate. What exactly do you store?
Minute 10:20, rolled up
A roll-up of a counter uses the gaps between consecutive samples: minute 10:20's increase is taken over the six gaps from 10:20:00 to 10:21:00. A gauge's minute holds the six samples 10:20:00, :10, …, :50. The roll-up runs 30 seconds after the minute ends (at M + 1 min 30 s), so the 10:21:00 sample has landed.
Synthesizing vector architecture diagram...
What to notice: the panel "Raw, minute 10:20" holds 7 samples per counter and 6 per gauge; the panel "1-minute rows" holds one number per counter and four for the gauge. Both give the same answers, mean 0.127 s and a p99 estimate of 2.035 s, because every stored number can be added to other minutes' and other pods' numbers.
Events 9 to 11
| # | Time | Event |
|---|---|---|
| 9 | 10:21:30 | The roll-up of minute 10:20. p4: _count +120, le="1" +0, le="2.5" +120, _sum +240 (120 × 2.0). p1: _count +1,200, le="0.1" +1,080, _sum +78. 44 rows written: one per counter (28), and four per gauge (16) |
| 10 | same | The minute's mean latency from the rows: Σ_sum ÷ Σ_count = (3 × 78 + 240) ÷ (3 × 1,200 + 120) = 474 ÷ 3,720 = 0.127 s. Averaging the four pods' means instead: (3 × 0.065 + 2.0) ÷ 4 = 0.549 s |
| 11 | same | p4's gauge for minute 10:20: samples 0, 4, 4, 4, 4, 4, so count 6, sum 20, min 0, max 4 (2 a second × 2 s = 4 in flight once the incident is 2 s old). Minute 10:21: 4, 4, 4, 4, 0, 0 |
What merges, and what doesn't
A value merges if the value for a longer window, or for more pods, can be computed from the values of the parts. Only those belong in a roll-up.
| Kept per window | Merges by | A query uses it for |
|---|---|---|
count (gauge samples) | adding | the denominator of an average |
sum (gauge samples) | adding | average = Σsum ÷ Σcount |
min, max | min of the mins, max of the maxes | the slowest moment, a spike |
| A counter's increase | adding | rate = increase ÷ window; mean latency = Σ_sum increase ÷ Σ_count increase |
| A bucket's increase | adding | percentiles, computed last from the summed buckets (Part 6) |
first, last | first of the firsts, last of the lasts | the open and close of a price bar |
| An average without its count | no: 0.549 s vs 0.127 s | |
| A percentile | no (Part 6) | |
| A distinct count ("unique users") | no; a HyperLogLog sketch does (Part 6) |
A stock chart's one-minute bars (open, high, low, close, volume) are a roll-up of this kind: first, max, min, last and sum, all mergeable, so a 5-minute bar is built from five 1-minute bars (the mobile stock-trading loop, R1.4). The payment loop's revenue account is summed "into per-minute totals" by a roll-up job in the same way (step 3.1 of the payment loop).
The type decides
A counter needs one value per window, its increase: the mean, the rate and every percentile come from increases. A gauge needs four: count, sum, min and max. A store that doesn't know the type must keep all five for every series: Thanos's downsampled chunks hold count, sum, min, max and counter for each one. (A series in delta form, where each sample is already the increase since the last, has no resets to handle, but a lost sample is a lost increase.)
Counters and restarts
p2's process restarts at 10:40:03 (event 13). Its counters start again from 0, and the 60 requests it counted between its 10:40:00 scrape and the restart are lost with the old process.
Synthesizing vector architecture diagram...
What to notice: the counter climbs to 120,000 at 10:40:00, then drops to 140 at 10:40:10 because the new process counts from 0. Taken naively, minute 10:40's last value minus its first is 940 − 120,000 = −119,060. The reset-adjusted increase treats the drop as a restart: 140 + 5 × 200 = 1,140.
The roll-up walks the gaps and treats any drop as a restart from 0:
textINCREASE(one counter's samples, window (M, M + w]) total = 0 ; gaps = 0 for each pair of consecutive samples (t0, v0), (t1, v1) with t1 in (M, M + w]: if v1 >= v0: total = total + (v1 - v0) else: total = total + v1 the counter restarted from 0 gaps = gaps + 1 if gaps == 0: write no row no data is not zero return total
Snapshot T3, 10:41:30: p2's _count across the restart (event 13)
| Value | |
|---|---|
| Samples 10:40:00 → 10:41:00 | 120,000 → 140 → 340 → 540 → 740 → 940 → 1,140 |
last − first inside [10:40, 10:41) | 940 − 120,000 = −119,060 |
| Reset-adjusted increase | 140 + 5 × 200 = 1,140 |
Requests p2 served in the minute | 1,200: the 60 counted by the old process after its last scrape are gone |
Only a counter can lose these increments, and only a few seconds' worth per restart. The query side of the same rule is Part 8's rate().
Keeping the last point loses the spike
Side row W replays the two hours with a cheaper roll-up, decimation: keep the last sample of each window, which is what VictoriaMetrics' downsampling does ("leaving the last sample per each interval"). It is compared with D4, the mergeable roll-up, on the same data, from both sides.
| Question | Decimation (last sample per window) | Mergeable roll-up (D4) |
|---|---|---|
p4's gauge, minute 10:20 | 4 (its 10:20:50 sample) | max 4 |
p4's gauge, minute 10:21 | 0 (10:21:50) | max 4 |
p4's gauge, hour 10 | 0 (10:59:50): the spike is gone | max 4 |
| Hour 10's mean latency | 0.065775 s (a 1-hour tier decimated straight from raw keeps only 10:59:50, so p2's restart falls inside the hour and it counts 174,940 requests instead of 223,140; its mean and p99 barely move) | 0.065774 s |
| Hour 10's p99 estimate | 0.462 s | 0.462 s |
p2's _count, minute 10:40 (a restart inside it) | 940 | 1,140 |
| Work and size | no job logic; one value per window | a job that knows each type; four values per gauge window |
The last value of a cumulative counter still carries every increase before it, so decimated histograms give the same mean and p99 as D4. What decimation loses is a gauge's spikes between the kept samples, and a counter's reset inside a window, where the kept samples no longer show the drop. Side row G in Part 14 goes one step further: a pod that freezes for 3 seconds between two scrapes shows nothing at all in its gauge, however it is rolled up, while its histogram counts every slow request. A gauge can't see between its samples; a counter can.
The hour from 60 minutes
Event 12. At 11:05:00 hour 10 is rolled up from its 60 minute rows: sums of increases, sums of counts and sums, the min of the mins and the max of the maxes. Checkout's _count for the hour is 223,140 (62 × 3,600 = 223,200, less p2's 60 lost requests). p4's gauge for the hour is count 360, sum 36, min 0, max 4: a mean of 36 ÷ 360 = 0.1 in flight, and a maximum that still shows the incident. The reference run also rolled hour 10 straight from the raw samples: for all 32 series, the two agree exactly. That is the test of a mergeable roll-up.
Set, never add
A roll-up job crashes and is retried; a late sample makes a window run again (Part 9). Every run of a window must leave the same row. So each row is keyed by (series, tier, window start), the whole window's value is computed, and the row is set to it: a second run writes the same value again. If the job instead adds its value to a stored total, a second run doubles it. The ad-click loop learned this in a failure: "Overwrite, never add, when work can be retried" (R1.9 of the ad-click loop); the job scheduler's meter records are keyed so that a re-sent one "overwrites itself in the daily rollup instead of adding to it" (step 3.5 of the job scheduler loop). Part 11 runs the retry both ways.
Where an add can't be avoided, it must be applied exactly once. A stream processor that adds each closed minute into its hour relies on its checkpoints for that ("The roll-up keeps the sum; a closed minute is added once", Part 3 of Event Time, Watermarks & Checkpoints); a store that applies adds guards each one with the last window applied, as the URL shortener's last_window condition does (step 2.2 of the URL shortener). Absolute values in general are Idempotency & Effectively-Once Processing, Part 5 there.
The canary's 90 seconds are rolled into minute 10:20. What exactly do you store, so that a year from now you can still compute that minute's mean latency, its slowest moment and its p99?
What to remember from Part 5
- Store sums and counts, never averages; min and max, never a sample.
- Roll up a counter as its reset-adjusted increase.
- A roll-up row is set for its window, never added to.
Part 6. Percentiles don't roll up
The hook's dashboard said the incident's p99 was 0.8 s by averaging the four pods' p99s; the truth was 2.0 s. A per-minute recording rule (D3) did better: it computed each minute's p99 from all pods' buckets. But averaging its 60 values for hour 10 gives 0.503 s, for an hour whose p99 was 0.4 s. Both mistakes come from treating a percentile as if it merged.
Why a percentile can't be merged
A p99 is the latency that 99% of requests beat: a rank in a sorted list. It depends on every value in the list, so two p99s can't make a third. In minute 10:20, three pods' p99 is 0.4 s and the canary's is 2.0 s. The fleet's p99 depends on how many requests each pod served, and on where the rest of each pod's latencies lie: 120 of the minute's 3,720 requests (3.2%) took 2 s, so the fleet's p99 is 2.0 s. No average of 0.4, 0.4, 0.4 and 2.0 gives that, and no weighting either: weighted by requests it is (3 × 1,200 × 0.4 + 120 × 2.0) ÷ 3,720 = 0.45 s. The p99 values alone don't carry what the merge needs.
| p99 of … | Incident (minute 10:20) | Hour 10 |
|---|---|---|
| The truth: nearest rank over every request | 2.0 s | 0.4 s |
| Average of each pod's p99 | 0.8 s | 0.326 s (240 pod-minutes) |
| Average of each minute's fleet estimate (D3) | 0.503 s | |
Estimate from the summed buckets (histogram_quantile) | 2.035 s | 0.462 s |
| A log-bucketed sketch, 1% relative error (DDSketch) | 1.994 s | 0.402 s |
Sum the buckets, then estimate
A histogram's bucket counts do merge: the number of requests at or below 1 s across four pods and 60 minutes is the sum of the counts. So the rule is: sum each bucket across pods and windows, and compute the percentile last. Prometheus's histogram_quantile does the last step: find the bucket where the rank falls, and interpolate linearly inside it (the lowest bucket starts at 0).
Event 14, minute 10:20. Summed across the four pods, the buckets' increases are: ≤ 0.1 s: 3,240; ≤ 0.5 s: 3,600; ≤ 1 s: 3,600; ≤ 2.5 s: 3,720; +Inf: 3,720. The rank for p99 is 0.99 × 3,720 = 3,682.8. That is above the 3,600 at or below 1 s, so it falls in the bucket (1, 2.5], which holds 120 requests; it is (3,682.8 − 3,600) ÷ 120 = 0.69 of the way in. The estimate is 1 + 0.69 × 1.5 = 2.035 s: above the slowest real latency, 2.0 s, because interpolation assumes the 120 requests are spread evenly across the bucket.
Event 15, hour 10. From the 1-hour rows: ≤ 0.1 s: 201,366; ≤ 0.5 s: 222,960; ≤ 1 s: 222,960; ≤ 2.5 s: 223,140; +Inf: 223,140. The rank is 0.99 × 223,140 = 220,908.6, which falls in (0.1, 0.5], holding 21,594 requests: (220,908.6 − 201,366) ÷ 21,594 = 0.905 of the way in, so 0.1 + 0.905 × 0.4 = 0.462 s. The true p99 is 0.4 s: only 180 of 223,200 requests were slow (0.08%). The incident survives in the hour only as the (1, 2.5] bucket's 180 and the gauge's max of 4. The p99 of a day, a month or a year is computed the same way from the 1-hour tier's buckets, never from stored p99s.
Synthesizing vector architecture diagram...
What to notice: in the panel "Right path", only bucket counts move from box to box, and each step is an addition; the percentile appears once, at the end. In the panel "Wrong path", percentiles are computed first and then averaged, which gives 0.8 s for an incident whose p99 was 2.0 s, and 0.503 s for an hour whose p99 was 0.4 s.
The bucket decides the error
A bucket estimate is only known to lie inside its bucket. 2.035 s could have been anything from 1 s to 2.5 s; 0.462 s anything from 0.1 s to 0.5 s. Two rules follow:
- Put bounds near the targets you alert on. If the SLO is "p99 under 300 ms", a bound at 0.3 turns a 0.1-to-0.5 guess into a yes or no.
- Put the top finite bound above your worst case. In the
+Infbucket, Prometheus returns "the upper bound of the second highest bucket".
Side row B (a copy): the bounds end at 1 s (0.1, 0.5, 1, +Inf). Minute 10:20's rank now falls in the +Inf bucket, and the estimate is 1.0 s, the highest finite bound. Part 10's alert "p99 above 1 s" can never fire: the answer is capped at exactly 1.
The gaming leaderboard loop (step 3.2) applies the same rule on purpose: 10,000 buckets at last season's quantiles, rank by the counts above plus interpolation inside the player's bucket, and "The error is at most the population of the player's bucket".
Sketches with a guarantee
A fixed bucket's error is absolute and can be huge. A log-bucketed histogram makes the error relative instead: its bucket bounds grow by a fixed ratio, so every bucket is, say, 2% wide wherever it sits. DDSketch is the clearest version. With a relative accuracy α = 1%, bucket i holds the values in (γ^(i−1), γ^i], where γ = (1 + α) ÷ (1 − α) ≈ 1.0202. A quantile is answered with the bucket's representative value, 2γ^i ÷ (γ + 1), which lies within 1% of every value in the bucket. Sketches merge by adding bucket counts, exactly like a histogram. Its paper calls it "the first fully-mergeable, relative-error quantile sketching algorithm with formal guarantees".
Event 16 (design D5). Each pod exports its latencies in DDSketch buckets instead of five fixed ones. The merged sketch answers 1.994 s for minutes 10:20 and 10:21 (true 2.0 s, 0.32% off) and 0.402 s for hour 10 (true 0.4 s, 0.62% off). Each pod-minute occupies 3 buckets on p1 to p3 (0.04, 0.18 and 0.4 s), 1 on p4 (2 in minute 10:21, when both 0.04 and 2.0 appear), and 4 in the merged fleet sketch. Prometheus's native histograms and OpenTelemetry's exponential histograms belong to this family: in Prometheus, "each bucket boundary is the previous boundary times 2^(2^-n)".
Two other sketches come up in interviews. t-digest "can be merged" and keeps its accuracy at the tails ("part per million accuracy for extreme quantiles"), with its error defined by rank rather than value. HDR histograms fix their precision as "the number of significant digits" across a stated range.
Histogram, sketch or summary, on equal terms
| Fixed-bucket histogram | Log-bucketed histogram (DDSketch, native, exponential) | Summary (finished quantiles from each client) | |
|---|---|---|---|
| Merges across pods and time | Yes: add the counts | Yes: add the counts | No |
| Error | Up to the bucket's width (2.035 s for a 2.0 s truth) | Within α of the true value (1%) | Within a configured error (Prometheus: "generally very low"), but for one process and its sliding window only |
| Size | One counter per bound (5 here) | One counter per occupied bucket: 3 per pod-minute here, a few hundred at worst | One value per quantile |
| Chosen in advance | Every bound | Only the accuracy α | The quantiles and their window |
| Client cost | One increment | A logarithm and one increment | A streaming quantile structure per process |
| Price on Managed Prometheus | One sample per bucket series | A native histogram is metered at 0.25 of a sample per populated bucket | One sample per quantile series |
A P99 alarm across many hosts must therefore be computed from merged buckets or sketches, never from the hosts' own p99s. The same goes for per-install summaries in the mobile paging loop (step 3.6: "counts and percentiles" per install per day): a fleet-wide p95 needs each install's histogram, with the same bounds everywhere.
Distinct counts
"Unique users per day" doesn't add up into a week either: a user active on Monday and Tuesday would be counted twice. A HyperLogLog sketch per day does merge, by taking the maximum of each register: step 3.5 of the URL shortener keeps one per link per day, about 12 KB, with a standard error of 0.81%, and says it plainly: "sketches merge". HyperLogLog's internals are outside this page; Bloom Filters & Counting Filters covers the family of sketches it belongs to.
You have each minute's p99 for the last hour. What is the hour's p99? What would you have needed to keep instead?
What to remember from Part 6
- Never average percentiles; merge buckets or sketches, then compute.
- A bucket estimate is only known to lie inside its bucket: pick bounds near your targets, and a top bound above your worst case.
- Distinct counts merge only as sketches.
Part 7. Tiers: how long each resolution lives
Kept raw at 16 bytes a sample, the metrics loop's Round 2 platform (100 million series, 10 million samples a second) would write 10 million × 86,400 × 16 B = 13.8 TB a day. Even compressed to the paper's 1.37 B, five years of raw 10-second data is 1.18 TB × 1,826 days ≈ 2.2 PB (R2.2 of the metrics loop), and a five-year graph of one series would decode 5 × 365 × 8,640 ≈ 15.8 million points. Tiers fix both problems. Here is what each one costs, counted honestly.
What each tier keeps
| Tier | Ours | Metrics loop | CloudWatch | Thanos (its compactor) | Managed Prometheus |
|---|---|---|---|---|---|
| Raw | 10 s, 15 days | 10 s, 30 days | Sub-minute points for 3 hours; 1-minute points for 15 days | Raw, kept as long as you set | Raw only: 150 days by default, up to 1,095 |
| Middle | 1 minute, 90 days | 5 minutes, 180 days | 5 minutes, 63 days | 5 minutes, for blocks older than 40 hours | none described |
| Top | 1 hour, 2 years | 1 hour, 5 years | 1 hour, 455 days | 1 hour, for blocks older than 10 days | none described |
Counting values, not windows
A roll-up doesn't store "one row per window". It stores values: one per counter per window, four per gauge per window (Part 5). So count values.
Event 18: values kept a day
| Raw | 1-minute tier | 1-hour tier | |
|---|---|---|---|
| One counter | 8,640 samples | 1,440 increases | 24 |
| One gauge | 8,640 samples | 1,440 × 4 = 5,760 | 24 × 4 = 96 |
| Our 32 series (28 counters, 4 gauges) | 32 × 8,640 = 276,480 | 28 × 1,440 + 4 × 5,760 = 63,360 | 28 × 24 + 4 × 96 = 1,056 |
A tier's size is then one product:
For bytes a value we plan with 2 bytes for every tier: the metrics loop's figure, and Managed Prometheus's own pricing example, which estimates storage as "2 bytes * number of metric samples". The reference run also measured our data: raw samples at 2.41 bytes (Part 4), and roll-up values at about 0.36 bytes. That second figure is an artifact of perfectly regular traffic, where every window's increase is the same number and XORs to zero; real traffic won't give it, so we don't plan with it.
Event 19: stored at retention, 2 bytes a value
| Tier | Values a day | Days kept | Stored |
|---|---|---|---|
| Raw | 276,480 | 15 | 276,480 × 15 × 2 B = 8.3 MB |
| 1-minute | 63,360 | 90 | 63,360 × 90 × 2 B = 11.4 MB |
| 1-hour | 1,056 | 730 | 1,056 × 730 × 2 B = 1.5 MB |
The 1-minute tier is the biggest: it is only 4.4 times smaller a day than raw, and it is kept 6 times longer. Roll-ups exist to answer long ranges fast and to outlive the raw data, not mainly to save space. Thanos's documentation says so plainly: keeping every resolution "doesn't save you any space", and "The goal of downsampling is to provide an opportunity to get fast results for range queries of big time intervals".
Synthesizing vector architecture diagram...
What to notice: the raw line climbs for 15 days and levels off at 4.15 million values; the 1-minute line keeps climbing until day 90 and levels off higher, at 5.70 million; the 1-hour line is still rising slowly at two years, at 0.77 million. The middle tier holds the most. (The axis is not evenly spaced.)
The "60× smaller" claim
"Rolling up to 1 minute cuts storage 60×" is true only for 1-second data kept as one value a window.
Event 20
| Roll-up | Values before | Values after | Smaller by |
|---|---|---|---|
| A gauge, 10 s to 1 minute | 6 samples | 4 aggregates | 6 ÷ 4 = 1.5× |
| A counter, 10 s to 1 minute | 6 samples | 1 increase | 6× |
| 1 s to 1 minute, one value | 60 samples | 1 | 60× |
| Our mix, per day | 276,480 | 63,360 | 4.4× |
This is why the metrics loop skips a 1-minute tier and keeps 5-minute and 1-hour tiers instead: from 10-second data, a four-aggregate 1-minute tier saves little and is kept for months.
Hot, warm and cold
Tiers of resolution are one axis; tiers of storage are another. The newest hours live in the head, in memory, because alerts and dashboards read them constantly. Recent blocks sit on the ingesters' local disks; older blocks move to object storage, which is cheap and slower to read the first time (R2.5 of the metrics loop). Managed stores draw the same line: Timestream for LiveAnalytics has "a memory store for recent data and a magnetic store for historical data", and OpenSearch Service moves indexes from hot storage to UltraWarm and then to cold. Part 13 has the details.
Six designs on equal terms
Side row X replays the same two hours under six designs and asks each the same questions. Each design gets its due: decimation needs no job logic and keeps one value a window; finished statistics are the smallest and ready to draw; raw forever keeps every question exact. The last column, "an hour from its minute rows", stands for re-rolling after late data: every design can be recomputed from raw while raw exists (15 days), but after that only a design whose rows merge can rebuild a coarser window.
| Design | What an old window keeps | Bytes a day (32 series, 2 B a value) | Hour 10's mean | Hour 10's p99 | p4's gauge spike | A reset inside a window | An hour from its minute rows |
|---|---|---|---|---|---|---|---|
| D0 a row per sample | every sample, 50 to 100 B a row | 13.8 to 27.6 MB | exact | exact | kept | handled | only from raw |
| D1 raw forever | every compressed sample | 0.55 MB, never deleted: 404 MB after two years | exact | exact | kept | handled | only from raw; 15.8 million points per series for five years |
| D2 decimation | the last sample of each window | 0.092 MB | 0.065775 s, right | 0.462 s, as D4 | lost (0) | lost (940, not 1,140) | yes: the last of the lasts |
| D3 finished statistics | each minute's fleet mean and p99, each pod's mean in flight | 0.017 MB | 0.065774 s: right here, because each minute's mean came from sums and the minutes carried nearly equal counts; minutes of unequal traffic would skew it, and a fleet number can't be split by pod | wrong: 0.503 s | lost: mean in flight 0.1 | not kept | no: averages of averages |
| D4 mergeable roll-ups | count, sum, min, max; increases | 0.127 MB | 0.065774 s | 0.462 s estimate, true 0.4 | kept (4) | handled (1,140) | yes, equal to the hour from raw |
| D5 D4 with log buckets | as D4, with 1% log buckets instead of 5 fixed ones | 0.098 MB | 0.065774 s | 0.402 s | kept | handled | yes |
D5 comes out smaller than D4 here only because our latencies take four distinct values, so each pod-minute occupies 1 to 3 log buckets against D4's five fixed ones; real latencies spread over more buckets.
Synthesizing vector architecture diagram...
What to notice: "D0" and "D1" answer everything but cost the most or grow forever. "D2" fails only on the gauge's spike and a reset inside a window. "D3" is the smallest and gets the p99 wrong; its mean is right only because every minute had about the same traffic, and it can't be split by pod. "D4" and "D5" answer every question from a small, mergeable roll-up.
Raw for 15 days, 1-minute for 90, 1-hour for 2 years. Which tier costs the most, and what would a 1-second raw tier with one aggregate have saved?
What to remember from Part 7
- Count the values a tier keeps: aggregates per window × windows × days kept.
- A 1-minute tier over 10-second data saves 1.5× for a gauge, 6× for a counter, not 60×.
- Tiers exist to answer long ranges fast and to outlive raw data.
Part 8. Reading it back: tiers, steps, rates and gaps
At 12:10:00 someone opens checkout's dashboard, with three panels over three ranges. For the longer ones, imagine the tiers have aged for a year (the arithmetic doesn't change). Each panel is a range query: a start, an end and a step, evaluated at start, start + step, …, end, both ends included. Which tier answers each, and how many values does it read?
Three panels, three tiers
Event 21
| Panel | Step | Points | Tier read | Values read per series | Raw would read |
|---|---|---|---|---|---|
| Last 2 hours | 10 s | 7,200 ÷ 10 + 1 = 721 | raw | 725 samples, 10:09:10 to 12:09:50 (a 1-minute rate window reaches 5 samples back; the 12:10:00 sample hasn't landed yet) | the same |
| Last 7 days | 5 minutes | 604,800 ÷ 300 + 1 = 2,017 | 1-minute: each point sums 5 minute rows | 7 × 1,440 = 10,080 rows (the newest point's 5 minutes are read raw, 30 samples) | 60,480 samples |
| Last year | 1 day | 366 | 1-hour: each point sums 24 hour rows per bucket, then takes the quantile | 365 × 24 + 24 = 8,784 rows | 3,153,600 samples |
The year panel's p99 is written histogram_quantile(0.99, sum by (le) (increase(checkout_duration_seconds_bucket[1d]))): at each daily point, the increase of each bucket over the day comes from 24 hourly rows, the four pods are summed per bucket, and the quantile comes last (Part 6).
A query reads the coarsest tier whose resolution is no coarser than its step, whose retention covers the range, and whose rows have been written. It checks that last condition per point:
textCHOOSE_TIER(step, the point's range) for tier in [1 hour, 1 minute, raw]: coarsest first if tier.resolution > step: skip if the range is older than tier.retention: skip if any row of the tier for the range isn't written yet: skip return tier raw qualifies within its 15 days at a step coarser than the tier, aggregate its rows over each step: sum the increases, max of the maxes, sum the buckets and then take the quantile report the tier used the loop's resolution_used field
So at 12:10:00 the 7-day panel's newest point, which covers minutes 12:05 to 12:09, is read from raw: minute 12:09 isn't rolled up until 12:10:30. Every other point comes from 1-minute rows. The response names the tier it used, so a graph can label itself "1-minute data" (the loop's resolution_used, R2.3).
Snapshot T5, 12:10:00: the three queries
| Query | Tier | Points | Values read per series |
|---|---|---|---|
| 2 h at 10 s | raw | 721 | 725 |
| 7 d at 5 min | 1-minute (newest point raw) | 2,017 | 10,080 rows + 30 raw samples |
| 1 y at 1 d | 1-hour | 366 | 8,784 |
Steps and alignment
Event 22. A panel opened at 12:00:07 asks for 10:00:07 to 12:00:07 at a 60 s step: 121 points, each at :07 past a minute. The next person to open it, at 12:00:31, asks a different question with different points, and neither matches anything cached. Aligned to the step, both ask for 10:00:00 to 12:00:00: the same 121 points, the same answer. So the query frontend aligns start and end to the step, splits long ranges into days, and caches the answers for past days. The cache key includes everything that changes the answer: tenant, query, step, range, tier, and a per-day generation that a re-roll or a backfill bumps (Part 9). The newest 10 minutes, still open to out-of-order samples, are never cached (step 3.4 of the metrics loop; caching in general is Caching & Invalidation).
Rates across a restart
Event 23. p2's counter restarted at 10:40:03. Prometheus's rate() adjusts resets per series ("Breaks in monotonicity (such as counter resets due to target restarts) are automatically adjusted for") and extrapolates to the ends of its window.
| Evaluated at | rate(p2 _count[1m]) | rate(p2 _count[5m]) | rate(sum(_count)), all pods | sum(rate(_count[1m])), all pods |
|---|---|---|---|---|
| 10:40:00 | 20.0 | 20.0 | 62.0 | 62.0 |
| 10:40:10 to 10:40:50 | 18.8 | 19.79 | 5,100.8 | 60.8 |
| 10:41:00 | 19.0 | 19.79 | 62.0 | 61.0 |
| 10:41:10 | 20.0 | 19.79 | 62.0 | 62.0 |
Per series, the rate dips just under 20 because of the 60 requests lost with the old process. The rate of the summed counter is absurd: the sum drops by 119,440 at 10:40:10 (from 372,000 to 252,560), rate() reads the drop as a reset of the whole sum, and adds the old value back, reporting about 5,100 requests a second for one window. Prometheus's documentation gives the rule: "always take a rate() first, then aggregate. Otherwise rate() cannot detect counter resets". The roll-up follows the same order: each series' increase is reset-adjusted first (Part 5), and only then summed.
A rate needs at least two samples in its window, so windows must be several scrape intervals long. rate(…[20s]) is fine while every scrape arrives, and returns nothing the moment one is missed (event 25 below).
Staleness and gaps
At 11:30:00 the canary pod is removed, and its target leaves service discovery (event 24). A query at an instant takes each series' newest sample within the last 5 minutes (the "lookback period is 5 minutes by default"). Without further help, p4 would be drawn at its last value for another 5 minutes. A pull system knows better: when a target disappears, the scraper writes a staleness marker into each of its series, and every evaluation at or after 11:30:00 finds no p4 at all.
Synthesizing vector architecture diagram...
What to notice: the store that receives the staleness marker drops p4 at 11:30:00, the moment its target leaves. The push store never hears that p4 left, so it keeps returning the 11:29:50 sample to every query less than 5 minutes after it, and p4 vanishes only at 11:34:50.
Events 24 and 25, and side row S
| # | Time | Event |
|---|---|---|
| 24 | 11:30:00 | p4's series end at 11:29:50 and carry a staleness marker at 11:30:00. Every evaluation at or after 11:30:00 returns no p4 |
| S | (side) | A push client without staleness markers (an OTLP- or StatsD-style exporter that just stops sending): every query less than 5 minutes after 11:29:50 still returns p4's last samples. Its gauge is drawn at its last value until just before 11:34:50 |
| 25 | (side) 11:40:10 | On a copy, p1's 11:40:10 scrape times out. rate(…[20s]) at 11:40:25 sees one sample in (11:40:05, 11:40:25] and returns nothing; rate(…[1m]) still has 5 samples and answers 20.0 |
A scraper that remote-writes (including Prometheus in agent mode, which keeps "the same scraping APIs, semantics") sends its markers along: the remote-write specification says senders "MUST send stale markers when a time series will no longer be appended to". A push pipeline without them draws ghosts for up to the lookback period.
No data is not zero. A gap must be drawn as a gap. Some tools fill it: CloudWatch's metric math AVG over an array of series says "Missing values are treated as 0", which drags an average down. And a rule that checks "p99 above 1 s" is silent both when things are fine and when there is no data at all, so pair it with an explicit rule for absence (absent(), R1.9 of the metrics loop).
One more reading hazard: a host whose clock runs behind makes its newest samples look late, so a rule over the last minute sees a dip. The metrics loop evaluates rules 30 seconds in the past (query_offset, R2.8) to leave room for it.
A 7-day panel at a 5-minute step, a year at 1 day, and the last 2 hours at 10 s: which tier answers each, and how many points come back?
What to remember from Part 8
- Each query reads the coarsest tier no coarser than its step.
- Take rates per series first, then sum; rates handle resets, sums don't.
- A stale or missing series is absent, never zero.
Part 9. Late data after a roll-up
At 11:05:00 p3's node loses its network. Its agent keeps scraping every 10 seconds and buffers what it can't send. At 11:12:00 the network is back, and at 11:12:05 the agent sends its 42 buffered scrapes (11:05:00 to 11:11:50) in one request, oldest first: 336 samples. By then the roll-ups of minutes 11:04 to 11:10 have already run without them. What does the store do with them, and what does the roll-up do?
Late, or out of order?
Event 26. The store's newest time when the batch lands is 11:12:00 (the other pods' 11:12:00 samples arrived at 11:12:01). p3's own newest sample is 11:04:50. Each buffered sample is newer than its series' newest, so for its series it is in order, and all 336 are accepted and appended to their chunks. They would be accepted even with an out-of-order window of 0.
Late means "arrived long after its timestamp"; out of order means "older than a sample its series already has". A chunk can only be appended to (Part 4), so only the second is a problem for the store. The batch is late but in order.
The out-of-order window
Side row O (a copy): the agent sends the batch newest first. 11:11:50 arrives first and is in order; the other 41 are older than their series' newest sample. Whether they are kept depends on the store's out-of-order window: a sample older than its series' newest is accepted if its timestamp is within the window of the store's newest time ("as long as the timestamp of the sample is >= TSDB.MaxTime-out_of_order_time_window", in Prometheus's words), and goes into a small separate out-of-order chunk for its series, merged with the normal one when the block is written.
| Out-of-order window | Accepted if the sample is at or after | Rejected per series (of 41) |
|---|---|---|
| 0 (Prometheus's default) | nothing older than the series' newest | 41 |
| 60 s (Managed Prometheus's default for a new workspace) | 11:11:00: five more are kept (11:11:00 to 11:11:40) | 36 |
| 10 minutes (the metrics loop's choice) | 11:02:00 | 0 |
The whole rule, as the store applies it to each sample of a request:
textACCEPT(sample of series s at time t) head_max: the head's newest time when the request arrived min_valid = max(head_max - 1 hour, end of the newest written block) if t >= min_valid and t > newest(s): append to s's chunk in order elif window > 0 and t >= head_max - window: append to s's out-of-order chunk merged at block write else: reject; route it to the backfill path
The one-hour floor is why Part 3's head keeps three hours before it cuts a block: anything in the newest hour can still be appended in order. A window costs memory for out-of-order chunks, and only series that actually receive late samples have one; Managed Prometheus also caps out-of-order ingestion at 5% of a workspace's ingestion rate. Data older than any window takes a slower path: the backfill path of step 2.5 of the metrics loop, where an hourly job builds normal blocks for the past windows and the compactor merges them with the existing ones (vertical compaction). Side row V in Part 14 runs it: a partition that lasts 95 minutes.
Dirty windows and re-rolls
The samples are in the store now, but the 1-minute rows for 11:04 to 11:10 were written without them. A sample accepted after its minute was rolled up makes that minute dirty: the minute of the gap it ends, and the minute it starts. The next run re-rolls every dirty minute from raw, with set (Part 5), and re-rolls any hour already built from them.
Events 27 and 28
| # | Time | Event |
|---|---|---|
| 27 | 11:12:30 | The regular run rolls up minute 11:11 (complete: the batch landed at 11:12:05) and re-rolls 7 dirty minutes, 11:04 to 11:10. Minute 11:04 had missed only its last gap (11:04:50 → 11:05:00 belongs to minute 11:04), so p3's _count row goes from 1,000 to 1,200, a gain of 200. Minutes 11:05 to 11:10 had no p3 rows at all (no gaps, so no data, not zero); each now gets 1,200 |
| 28 | 12:05:00 | Hour 11 is rolled up from its 60 minute rows: _count 219,580 (62 × 1,800 + 60 × 1,800 − 20: p4 served until 11:30:00, but the 20 requests after its last scrape, 11:29:50, were never counted). The hour doesn't wait for p4's rows after 11:30, because p4's series are stale (Part 11). With an ignore policy (no re-roll) and an hour job that rolls up whatever minute rows exist, it would be 212,180, missing 200 + 6 × 1,200 = 7,400 of p3's requests, and every mean, rate and p99 for 11:04 to 11:10 would be computed without most of p3. With the completeness gate on, the hour would instead wait forever: 48 of p3's minute rows are never written |
Snapshot T4, 11:12:30: checkout's _count per minute, all pods
| Minute | As first rolled up | After the re-roll |
|---|---|---|
| 11:04 | 3,520 | 3,720 |
| 11:05 to 11:10 | 2,520 each | 3,720 each |
| Hour 11 (at 12:05) | 212,180 if ignored (no gate) | 219,580 |
Synthesizing vector architecture diagram...
What to notice: the panel "Rolled up 11:05:30 to 11:11:30" shows the minutes as first written, without p3. The batch in "11:12:05: batch lands" marks 7 minutes dirty, and "Re-rolled 11:12:30" rewrites them. In "Hour 11 at 12:05", the hour built after the re-roll counts 219,580; the dotted path, which never re-rolls and has no completeness gate, gives 212,180.
Every answer computed from those minutes before 11:12:30 was wrong and is now different. So a re-roll must also tell the results cache: it bumps the generation of the day it changed (Part 8), and a day's answers aren't cached until it is past the out-of-order window.
Early and often, or late and rarely
There are two ways to schedule roll-ups, and both are in production.
| Roll up early, re-roll what gets dirty (our main line) | Roll up late, once (Thanos: 5-minute roll-ups of blocks older than 40 hours) | |
|---|---|---|
| How fresh the tier is | 30 s after a minute ends (90 s after it begins) | 40 hours |
| Late data | Re-roll each dirty window, and bump the cache generation | Usually landed before the roll-up; our 7 minutes would need nothing |
| Work | Proportional to late data | One pass per block |
| Machinery | Dirty tracking, and set-by-key writes | None beyond the compactor |
| What readers must know | Recent windows are provisional; they may change until late data stops | A long query falls back to raw until the roll-up exists |
The loops use the first. The S3-like storage loop meters usage into hourly roll-ups, and "A batch that arrives late … marks that hour dirty, and the rollup is recomputed. Hours are final after a day." The ad-click loop recomputes the last three event hours on every run and overwrites them (step 1.5). CloudWatch is honest about its own open periods: "data aggregated between 7:00pm and 8:00pm begins to be visible at 7:00pm, then the values of that aggregated data may change as CloudWatch collects more samples during the period"; and it accepts time stamps "up to two weeks in the past". Whichever you choose, decide it once and say it: which windows are provisional, and when they become final. A streaming job's late path, its watermark and its acceptance horizon are Event Time, Watermarks & Checkpoints, Part 4 there.
Seven minutes of one pod's samples arrive at 11:12:05. The minutes they belong to were rolled up at 11:05:30 to 11:11:30. What do you do?
What to remember from Part 9
- Late but in order for its series is fine; out of order needs a window; older still needs backfill.
- Every roll-up that covered the window is dirty until re-rolled; so is the hour built from it.
- Decide once: re-roll, or mark the window as final and say so.
Part 10. Alerting on coarse data
The rule is "checkout p99 above 1 s", evaluated every 15 seconds. The canary's p99 was 2 seconds for 90 seconds. When can it fire? It depends entirely on which data the rule reads.
Event 30: three tiers, three fire times
| Rule reads | Fires | Why |
|---|---|---|
Raw data: histogram_quantile(0.99, sum by (le) (increase(checkout_duration_seconds_bucket[1m]))) > 1 | 10:20:30, 30 s after the incident starts; clears at 10:22:15 | The first evaluation whose 1-minute window holds more than 1% slow requests |
| The 1-minute tier: the same estimate from each minute's row, evaluated when the row is written | 10:21:30, as the incident ends; clears at 10:23:30 | Minute 10:20's row (estimate 2.035 s) is written at 10:21:30 |
| The 1-hour tier | never | Hour 10's estimate is 0.462 s |
Synthesizing vector architecture diagram...
What to notice: the raw rule fires 30 seconds into the 90-second incident. The 1-minute rule fires only when minute 10:20's row is written, at 10:21:30, the moment the incident ends. The hourly rule's bar stretches to 11:05 and ends with no alert: the hour's estimate, 0.462 s, never crosses 1 s.
What sets the delay
On a tier, an alert can't fire before the window closes and its roll-up runs:
For minute 10:20: 10:21:00 + 30 s + 0 (the rule runs as the row is written) = 10:21:30. For an hourly tier, at best 11:00:00 + 5 minutes, and only if the incident moves the hour's p99, which a 90-second incident can't.
On raw data the window slides, so the question is when the bad share inside it crosses the threshold. The p99 is above 1 s once more than 1% of the window's requests are slow: 0.99 × 3,720 = 3,682.8, so at least 38 of a minute's 3,720. The canary adds 2 slow requests a second, so about 19 seconds of them. The scrapes see them in steps of 20: 20 at the 10:20:10 scrape, 40 at 10:20:20, which lands at 10:20:21. The evaluation at 10:20:15 sees 20 slow of 3,100 in its window (0.65%): not yet. The evaluation at 10:20:30 sees 40 of 2,480 (1.6%): it fires. So detection on raw data is the time for the bad share to cross the threshold, plus the scrape interval and the landing delay, plus up to one evaluation interval.
| Evaluated at | Slow share in the window | p99 estimate | State |
|---|---|---|---|
| 10:20:00 | 0% | 0.459 s | ok |
| 10:20:15 | 0.65% | 0.485 s | ok |
| 10:20:30 | 1.61% | 1.570 s | firing |
| 10:21:00 to 10:21:30 | 3.23% | 2.035 s | firing |
| 10:22:00 | 1.61% | 1.570 s | firing |
| 10:22:15 | 0.65% | 0.485 s | ok |
A rule's for duration adds its own delay on top: the metrics loop's example rule waits for: 5m, so this 90-second incident, firing from 10:20:30 to 10:22:00, would never page at all. For a canary, that may be exactly what you want; choose it per rule (step 1.4 of the loop).
The hour dilutes
Event 31. The 1-hour tier doesn't just see the incident late; it barely sees it. Hour 10's mean latency is (3,600 × 3.98 − 180 × 0.04 + 180 × 2.0) ÷ 223,200 = 0.0658 s, against 3.98 ÷ 62 = 0.0642 s in a normal hour (3.98 s is one second's 62 latencies added up). The 180 slow requests are 0.08% of the hour. A coarse tier is for history, not for detection.
So: alert on raw data, or on recording rules computed from raw data every evaluation. Long windows, such as a 30-day SLO, come from recorded short sums added up, not from averaging a tier: the metrics loop builds its 3-day burn-rate ratio "from recorded 5-minute sums rather than raw samples" (step 3.3). The same loop budgets the whole path from scrape to page (R1.7). And a threshold rule is silent when data stops, so give it an absent() companion (Part 8).
On Managed Prometheus the minimum rule evaluation interval is 30 s by default (an adjustable quota), so our 15-second rule would evaluate half as often there.
The canary's p99 is 2 s for 90 seconds. When does "p99 above 1 s" fire on raw data, on the 1-minute tier and on the 1-hour tier?
What to remember from Part 10
- Resolution sets the earliest possible alert.
- Alert on raw data; roll-ups are for history.
- Long windows come from recorded sums, not from averaging a tier.
Part 11. When the pipeline breaks
Every Part so far assumed the roll-up jobs run on time, once, and the head never loses its memory. None of that holds for long in production. On copies of the main line, here is what goes wrong, and what keeps the tiers right anyway.
The job that ran twice
Side row R. At 11:21:30 the roll-up of minute 11:20 starts writing its 44 rows. After 22 (all of p1's and p2's: seven counters and four gauge values each), the worker crashes. The job is retried 10 seconds later and writes all 44.
| The retry writes with | Rows now wrong | Checkout's _count for 11:20 | Requests a second | p99 estimate | p1's gauge max |
|---|---|---|---|---|---|
| add | 22, each twice its value | 3,720 + 2,400 = 6,120 | 102 | 0.459 s | 4 (it was 2) |
| set | 0 | 3,720 | 62 | 0.459 s | 2 |
The doubled minute is hard to spot. p1's and p2's buckets, _sum and _count are all doubled together, so the p99 and the mean barely move (the mean goes from 0.0642 s to 0.0645 s). Only the request rate jumps, for one minute, and hour 11 now counts 2,400 requests that never happened.
Synthesizing vector architecture diagram...
What to notice: the retry writes all 44 rows, including the 22 the first attempt already wrote. With add, those 22 are counted twice; with set, the second write of each row leaves the same value, so the minute is right however many times the job runs.
A roll-up that is safe to run any number of times looks like this:
textROLL_UP(tier, window W) safe to run any number of times for each series s: if s is a counter: value = INCREASE(s's samples, W) Part 5 else: value = (count, sum, min, max) of s's samples in W if there is no data: write nothing missing, not zero else: SET row (s, tier, W) = value overwrite, never add mark W rolled up; clear W's dirty flag ROLL_UP_HOUR(H) if a series that isn't stale lacks a minute row in H: wait, and retry every minute else: SET each hour row = the merge of its 60 minute rows
A worker that fell behind
Side row K. The roll-up worker is down from 11:40:00 to 12:10:00. The runs due at 11:40:30 to 12:09:30 don't happen, so 30 minute windows (11:39 to 12:08) have no rows. Two things must not break:
- Queries at a 1-minute step find no rows for that range, so they fall back to raw data for those points (Part 8's rule, applied per point). They are slower, not wrong.
- Hour 11, due at 12:05, would be built from 39 of its 60 minutes. The hour job checks first: for every series that isn't stale, all 60 minute rows must exist. Minutes 11:39 to 11:59 are missing for
p1,p2andp3, so it waits.p4's missing rows after 11:30 don't count: its series have been stale since then, and waiting for them would block the hour forever.
Snapshot T6, copy K
| Time | Minute rows for hour 11 | Hour 11 | A 1-hour panel at a 1-minute step |
|---|---|---|---|
| 12:00:00 | 11:00 to 11:38 written; 11:39 onward missing | not yet due | 40 points from 1-minute rows, 21 from raw (11:40:00 to 12:00:00) |
| 12:05:00 | still 39 of 60 | due: waits (504 rows missing: 21 minutes × 24 series) | |
| 12:10:00 | the worker is back and rolls up all 30 missing minutes at once | rolled: 219,580, the same as the main line | all from 1-minute rows |
Run anyway at 12:05, hour 11 would have counted 39 minutes: 143,980 requests, missing 21 × 3,600 = 75,600 (60 a second, after the canary left). A coarse tier is complete only when all its inputs are: wait, or write it marked partial so that readers can tell. The worker's capacity is shared by the regular runs, the re-rolls of Part 9 and any catch-up, so a burst of late data or an outage delays every tier at once. The metrics loop's failure table calls this "The compactor falls behind" (R2.8) and alarms on the age of the oldest uncompacted block; the equivalent here is the age of the oldest missing roll-up.
An ingester that crashed
The head lives in memory. When the ingester holding it crashes, the head is rebuilt by replaying the WAL: on a copy (side row I, in Part 14), the ingester crashes at 11:50:00 and is back at 11:50:30, replaying the log's 31,688 samples and markers (and its 32 series entries) rebuilds all 32 series exactly, and the agents resend the 3 scrapes per series they couldn't deliver, in order. At the metrics loop's scale the replay takes much longer: about 33 minutes to rebuild a shard's head (R2.6), which is why that design keeps two replicas of every series so that queries and alerts are answered meanwhile (step 2.2). Durable logs and replay are Write-Ahead Log, fsync & Group Commit; a second copy is Replication, Quorums & Read-Your-Writes.
A corrupt block
Blocks are files, and files get damaged. The storage engine detects it with checksums and recovers by dropping the damaged range and fetching it again from another replica; the mechanics are in Part 10 of LSM-Trees & Compaction. And one failure needs no crash at all: a label that multiplies the series (Part 2).
The 11:20 roll-up job crashed halfway and was retried. What decides whether minute 11:20 is now right, or 22 of its rows are doubled?
What to remember from Part 11
- A roll-up must be safe to run twice: set by key.
- A coarse tier is complete only when all its inputs are: wait, or mark it partial.
- The head is rebuilt from its log; alerting needs a second copy meanwhile.
Part 12. End to end through the layers
One sample, many layers. At 10:20:10 p4's bucket le="2.5" reads 9,620: the first sample that shows the canary's slow requests. Here is everywhere it goes, and what each layer keeps.
Synthesizing vector architecture diagram...
What to notice: the path splits at the ingester's head. Alert rules read the head's raw data within seconds; history flows on to blocks in object storage and to the roll-up tiers, which the query frontend reads by step. Each box keeps its data for a different time, from the client's counters (until the process restarts) to the 1-hour tier (2 years).
Event 32: the 10:20:10 sample of p4's le="2.5" bucket
| When | Layer | What happens |
|---|---|---|
| 10:20:00.250 → 10:20:09.750 | Client library | 20 slow requests each increment the counters for le="2.5" and +Inf (and _sum, _count); nothing is sent yet |
| 10:20:10 | Agent | Scrapes p4: le="2.5" = 9,620, stamped 10:20:10.000 |
| 10:20:11 | Gateway, stream, ingester | The request passes the tenant's limits; the sample is appended to the WAL and to series 28's open chunk (a few bits: Part 4); queryable at once |
| 10:20:15 | Rule evaluator | Reads it: 20 slow of 3,100 in the window, not yet over 1% |
| 10:20:30 | Rule evaluator | With the 10:20:20 sample too: 1.6% slow, the alert fires (Part 10) |
| 10:21:30 | Roll-up | Minute 10:20's row for p4's le="2.5": increase 120 |
| 11:05:00 | Roll-up | Hour 10's row, merged from 60 minute rows |
| 13:00:11 | Ingester | The block 10:00 to 12:00 is written and uploaded to object storage; the WAL before 12:00 is dropped |
| 15 days later | Retention | The raw block is deleted whole (within about two hours of expiring) |
| A year later | Query frontend | "Yearly p99" reads hour 10's row from the 1-hour tier, summed with the day's other hours |
Layer by layer
| Layer | What it keeps | For how long | What it limits | When it fails |
|---|---|---|---|---|
| Client library | Counters and bucket counts, in process | Until the process restarts | Nothing | A restart resets its counters: resets are handled downstream (Part 5); a few seconds of increments are lost |
| Agent | A buffer of scraped samples | Minutes to hours (Prometheus's agent mode: "limited to a two-hour buffer") | Its own queue | Data is late but in order (Part 9); beyond its buffer, lost |
| Gateway | Nothing | Forbidden labels, label counts and lengths (400); samples a second per tenant (a token bucket, 429) | Its limits refuse a launch: they need owners | |
| Stream (the metrics loop's design) | Every write, in order | Hours | Throughput per shard | Consumer lag (Queues & Delivery Semantics) |
| Ingester head | WAL, open chunks, the label index | About 3 hours | Active series per tenant | Replay after a crash; a second replica answers meanwhile (Part 11) |
| Blocks in object storage | Immutable 2-hour blocks, merged | Raw retention | Slow first reads; cleanup lags | |
| Compactor and roll-ups | The tiers | 90 days, 2 years | Its own capacity, shared with re-rolls | Falls behind: queries fall back to raw, hours wait (Part 11) |
| Query frontend | Aligned, cached answers for past days | Until a generation bump | Samples per query | A cache key without the tier or generation serves stale answers (Part 8) |
| Rule evaluator and alert router | Alert state | Minutes | Rules per tenant; evaluation interval | Silent on missing data unless an absent() rule is paired (Part 8) |
What to remember from Part 12
- Each layer keeps the least it needs: counters in the client, bytes in the chunk, aggregates in the tiers.
- Limits sit where the cost is: series in the ingesters, samples at the gateway.
- Every tier and cache must know when late data changed its window.
Part 13. On AWS
AWS runs this mechanism in four managed services, and each makes a different choice about the three questions this page keeps asking: what is kept for an old window, how late a sample can be, and what you pay for (samples, series, or instances). Anything else, you run yourself.
Managed services that use it
| Service | What it provides | What AWS documents |
|---|---|---|
| Amazon CloudWatch (metrics) | A managed time-series store with fixed roll-up tiers, statistics per period, percentiles and metric math | Retention: "Data points with a period of less than 60 seconds are available for 3 hours"; 60-second points for 15 days, 5-minute for 63 days, 1-hour for 455 days (15 months); "After 15 days this data is still available, but is aggregated and is retrievable only with a resolution of 5 minutes. After 63 days, the data is further aggregated and is available with a resolution of 1 hour". A high-resolution custom metric is stored at 1 second (StorageResolution "Valid values are 1 and 60"). Identity: CloudWatch "treats each unique combination of dimensions as a separate metric, even if the metrics have the same metric name". Late data: "The time stamp can be up to two weeks in the past and up to two hours into the future"; an open period's values "may change as CloudWatch collects more samples during the period". Percentiles: "available for custom metrics as long as you publish the raw, unsummarized data points"; "CloudWatch needs raw data points to calculate percentiles", so a statistic set (Min, Max, Sum, SampleCount) supports them only when "The SampleCount value of the statistic set is 1 and Min, Max, and Sum are all equal"; "Percentile statistics are not available for metrics when any of the metric values are negative numbers". A client-side histogram is sent as Values and Counts arrays, "up to 150 unique values in each PutMetricData action". Metric math: RATE is "the difference between the latest data point value and the previous data point value, divided by the time difference in seconds", and in AVG over several series "Missing values are treated as 0". Price (US East): $0.30 a metric a month for the first 10,000, $0.10 for the next 240,000, $0.05 for the next 750,000, $0.02 above 1,000,000, "prorated by the hour" |
| Amazon Managed Service for Prometheus | A managed Prometheus-compatible store: raw samples, an out-of-order window, series limits, rule evaluation; no roll-up tiers described | "Metrics ingested into a workspace are stored for 150 days by default", configurable up to "1095 days (three years)". Active series per workspace: 50,000,000 by default, "up to a maximum of 1.5 billion"; "A series is active if a sample has been reported in the past 2 hours". Ingestion: 1,666,666 samples a second, "1/30 of the active series per workspace limit". "The default out-of-order time window when a new workspace is created is 60 seconds, and it can be configured up to a maximum of 600 seconds"; out-of-order ingestion is capped at 5% of the ingestion rate. Minimum rule evaluation interval: 30 s (adjustable). Queries: at most 95 days of "Query time range in days". Labels per series: 150. Price: $0.90 per 10 million samples for the first 2 billion (then $0.35 and $0.16), storage at $0.03 a GB-month estimated as "2 bytes * number of metric samples", queries at $0.10 per billion samples processed; a native histogram is metered at 0.25 of a sample per populated bucket |
| Amazon Timestream | Two products. For LiveAnalytics: a serverless store with a memory tier, a magnetic tier and scheduled queries for roll-ups. For InfluxDB: managed InfluxDB instances and clusters | LiveAnalytics: "we have made the decision to close new customer access to Amazon Timestream for LiveAnalytics, effective 6/20/25 … We recommend that new customers evaluate Amazon Timestream for InfluxDB". It keeps "a memory store for recent data and a magnetic store for historical data", moved "based upon user configurable policies" (memory retention 1 to 8,766 hours, default 6; magnetic 1 to 73,000 days). "Late-arriving data is data with a timestamp earlier than the current time and outside the memory store retention period. You must explicitly enable … magnetic store writes". Scheduled queries "compute aggregates, rollups, and other operations … and reliably writes the query results into a separate table", whose retention "is fully decoupled from that of source tables". InfluxDB: "the familiar open source version of InfluxDB on its 2.x branch", with a Multi-AZ standby; "High series cardinality is the primary driver of high memory usage" |
| Amazon OpenSearch Service | Index rollups, and hot, UltraWarm and cold storage, for time-stamped documents (logs and events) | "Index rollups … let you reduce storage costs by periodically rolling up old data into summarized indexes", keeping "only those fields aggregated into coarser time buckets", with the metrics "avg, sum, max, min, and value count": no percentiles. Upstream OpenSearch documents two more points that AWS's page doesn't: an avg under a date_histogram on a rollup index must be computed from sum and value_count (a rolled-up average is not stored finished), and a cardinality rollup metric uses HyperLogLog++. UltraWarm uses "Amazon S3 and a sophisticated caching solution", and warm indexes are "read-only"; Index State Management moves an index "from hot storage to UltraWarm, and eventually to cold storage. Then, it deletes the index" |
CloudWatch, Managed Prometheus, Timestream for InfluxDB and your own cluster, on equal terms
| CloudWatch | Managed Service for Prometheus | Timestream for InfluxDB | Your own (Thanos or Mimir on EKS, blocks in S3) | |
|---|---|---|---|---|
| Finest resolution | 1 s (high-resolution custom metrics); 60 s by default | Whatever you scrape (millisecond timestamps) | The engine's | Whatever you scrape |
| Roll-ups | Fixed: 5-minute after 15 days, 1-hour after 63 | None described: raw samples for the whole retention; build long views with recording rules | Whatever you build in the engine | Yours to schedule: Thanos 5-minute after 40 hours, 1-hour after 10 days |
| Retention | 3 h, 15 d, 63 d, 455 d by resolution | 150 days by default, up to 3 years | Your instance's storage | Yours, per resolution |
| Late data | Time stamps up to two weeks old | A 60 s window by default, up to 600 s | The engine's | Your out-of-order window and a backfill path |
| Percentiles | From raw values or Values/Counts arrays; statistic sets only when every value is equal | histogram_quantile from buckets, or native histograms | The engine's | histogram_quantile from buckets |
| You pay for | Each metric (each dimension combination) a month | Each sample ingested, GB stored and samples queried | Instances and storage | Servers, disks and S3 |
| Series limits | None stated; the price is per series | 50 million active series by default | Memory | Your ingesters' memory |
| Where it wins | No servers; built-in alarms; AWS's own metrics already there | PromQL, and 150 days of raw samples, with no servers | A managed InfluxDB when you want its query language | Cost at scale: the metrics loop's 100 million series cost about $76,000 a month to run itself, against about $425,000 a month for Managed Prometheus's ingestion alone |
Three limits are shared. A Managed Prometheus workspace's 50 million active series and 1,666,666 samples a second are shared by everything that writes to it, and out-of-order samples have their own 5% share. A query's sample budget (Managed Prometheus: 50 million samples per 24-hour interval of a query) is shared by every panel of a dashboard refresh. And CloudWatch's price per metric is multiplied by every dimension combination a service emits (Part 2's user_id: $6,000 a month).
Running it yourself
| Option | What it is | Facts and sizing |
|---|---|---|
| Prometheus on Amazon EC2 or Amazon EKS | The store of Parts 3 and 4: head, WAL, 2-hour blocks | "Ingested samples are grouped into two-hour blocks"; the WAL is written "in 128MB segments"; compaction merges blocks "up to 10% of the retention time, or 31 days, whichever is smaller"; retention defaults to 15 days; disk ≈ retention in seconds × samples a second × bytes a sample, at "an average of only 1-2 bytes per sample"; out_of_order_time_window defaults to 0; the lookback is 5 minutes. One server holds only what one server's memory holds (Part 2) |
| Thanos or Grafana Mimir on EKS, blocks in Amazon S3 | Prometheus scaled out: ingesters, blocks in object storage, a compactor, a query frontend | Thanos's compactor: "5m downsampling for blocks older than 40 hours", "1h downsampling for blocks older than 10 days", keeping count, sum, min, max and counter, with retention per resolution; keeping them all "doesn't save you any space". Mimir's sizing: "Memory: 2.5GB for every 300,000 series in memory", "CPU: 1 core for every 300,000 series in memory". The metrics loop's Round 2 is this design at 100 million series |
| VictoriaMetrics on EC2 or EKS | A Prometheus-compatible store | Its downsampling keeps "the last sample per each interval": decimation (Part 5). Right for counters and histograms, which are cumulative, but a gauge's spikes between kept samples and a reset inside an interval are lost |
| InfluxDB or TimescaleDB on EC2 | A time-series database, or PostgreSQL with time-series extensions | TimescaleDB's continuous aggregates are "refreshed in the background when new data is added", and real-time aggregates "combine pre-aggregated data with the most recent raw data" |
| Sizing in words | Memory from series (about 8.3 KB each, a planning figure), never from compressed bytes. Disk from the formula above. Object storage for blocks. Steady-CPU instances for ingesters. Alarms on active series per tenant, on refused samples, on the age of the oldest missing roll-up, and on WAL replay time |
Look-alikes that are not this mechanism
| Look-alike | Why it looks like this | Why it isn't |
|---|---|---|
| DynamoDB with TTL | "Rows per minute that expire" (the ad-click loop's minute rows) | A key-value store you shape into windows yourself. TTL deletes expired items within a few days, so readers must filter by age (the ad-click loop's R2.5 does); no chunk compression, no roll-up engine |
| Kinesis Data Streams or MSK retention | "Keeps 24 hours to a year of data" | Log retention of every record, in order; nothing is aggregated (Queues & Delivery Semantics) |
| S3 Lifecycle and storage classes | "Moves old data to cheaper tiers" | Moves or deletes whole objects by age; it never aggregates what is inside them |
| CloudWatch Metric Streams | "Streams our metrics out" | An export; each update carries "four default statistics; Minimum, Maximum, Sample Count and Sum", already rolled up per minute |
| CloudWatch Logs metric filters and the Embedded Metric Format | "Metrics from logs" | Ways to get data into CloudWatch metrics; the storage and roll-ups are CloudWatch's |
| Amazon Managed Service for Apache Flink | "It computes the minute and hour windows" | It computes windows over a stream and stores nothing itself; its output goes to a sink (Event Time, Watermarks & Checkpoints) |
| Amazon Managed Grafana | "Our metrics platform" | Dashboards only; it stores no samples |
| Athena over Parquet in S3 | "Hourly counts from the data lake" | A query engine over files: exact, slow, billed per byte scanned; no tiers or chunk compression of its own |
| ElastiCache sorted sets or counters | "Real-time counts and leaderboards" | In-memory structures (the leaderboard loop); no time-series retention or roll-up |
What to remember from Part 13
- CloudWatch rolls up on its own schedule (3 hours, 15 days, 63 days, 455 days) and prices per metric, so cardinality is money.
- Managed Prometheus keeps raw samples (150 days by default, up to 3 years) and describes no roll-up tiers; one query spans at most 95 days, so a year's panel is several queries.
- Timestream for LiveAnalytics is closed to new customers; Timestream for InfluxDB is the managed path.
Part 14. What you've learned
Back to the canary
The incident review asked for the 90 seconds' mean and p99. Averaging each pod's numbers said 0.549 s and 0.8 s; the truth was 0.127 s and 2.0 s. The capacity review asked for the yearly p99 with the incident on it. Averaging each minute's p99 said 0.503 s for hour 10; the hour's summed buckets give 0.462 s as an estimate, the truth is 0.4 s, and the incident survives as 180 requests in the (1, 2.5] bucket and a gauge maximum of 4. Here is what each piece did:
- Counting series, not samples (Parts 1 and 2) showed 32 series from two metrics, and why one
user_idlabel turns them into 280,000: 2.3 GB of memory and $6,000 a month, for the same traffic. - The head, the WAL and immutable blocks (Part 3) made every write an append and every deletion a whole block.
- Delta of deltas and XOR (Part 4) brought 16 bytes down to 2.41 on our series, from a quarter of a byte for a steady gauge to 7 for a float sum.
- Roll-ups that keep only what merges (Part 5): increases for counters,
count,sum,min,maxfor gauges, set by key, so minute 10:20 still answers 0.127 s andmax4 a year later. - Merging buckets before the percentile (Part 6) turned 0.8 s and 0.503 s into 2.035 s and 0.462 s, and a 1% sketch into 1.994 s and 0.402 s.
- Counting values per tier (Part 7) showed the 1-minute tier is the biggest, 11.4 MB, and that 60× is a myth at 10 s.
- Choosing the tier by step, rating before summing, and treating stale as absent (Part 8) kept the year panel at 8,784 rows, the rate at 19 instead of 5,100, and the removed canary off the graph.
- Re-rolling dirty windows (Part 9) put
p3's 7,400 late requests back into hour 11: 219,580, not 212,180. - Alerting on raw data (Part 10) fired at 10:20:30, while the 1-minute tier fired as the incident ended and the hour never did.
- Set-by-key writes and completeness gates (Part 11) kept a retried job from doubling 22 rows and a stalled worker from publishing an hour 75,600 requests short.
- Following one sample through every layer (Part 12) put each limit where its cost is.
The snapshots, side by side
| Snapshot | Where and when | What it showed |
|---|---|---|
| T1 | main, 10:00:10 | p4's eight series, two samples each, in their open chunks |
| T2 | main, 10:21:30 | Minute 10:20 as raw samples and as 44 rows, giving the same mean (0.127 s) and p99 estimate (2.035 s) |
| T3 | main, 10:41:30 | p2's counter across the restart: last − first = −119,060, the reset-adjusted increase 1,140, 60 requests lost |
| T4 | main, 11:12:30 | Minutes 11:04 to 11:10 re-rolled: 3,520 and 2,520 become 3,720; hour 11 is 219,580, not 212,180 |
| T5 | main, 12:10:00 | Three panels: raw 721 points, the 1-minute tier 2,017, the 1-hour tier 366 |
| T6 | copy K, 12:00 to 12:10 | 30 missing minutes, 21 panel points read from raw, hour 11 waiting until 12:10:00 |
What it costs
- Roll-up jobs that must keep up, be safe to run twice, and track which windows late data made dirty.
- Percentiles only as exact as the buckets, or a sketch format that every layer supports.
- Detection delay on coarse data: alerts must read raw data or recording rules built from it.
- Late data that forces re-rolls and cache invalidation, and recent windows that stay provisional until they are declared final.
- Kilobytes of memory per series, and a limit to enforce on every label.
- Storage that tiers add: kept together, the tiers are bigger than raw alone, and the middle one is often the biggest.
The whole story, event by event
Side rows run on copies of the main line and are marked "(side)".
| # | Time | Where | Event |
|---|---|---|---|
| 1 | 10:00:00 | main | The first scrape of the block: 32 samples. p1's _count 72,000, le="0.1" 64,800; p4's _count 7,200 |
| 2 | 10:00:10 | main | Each counter grows by a fixed step (p1's _count +200, _sum +13.0); the gauge reads 2, while the average in flight is 1.3 |
| 3 | text | main | The index: pod="p4", le="2.5" and the metric name intersect to series 28; memory 32 × 8.3 KB ≈ 266 KB, index 6.4 KB |
| C | (side) from 10:50 | copy | A user_id label: 280,000 series, 2.3 GB of head memory, about 28,000 samples a second; as CloudWatch metrics 40,000 pod-user histograms = $6,000 a month instead of $1.20 |
| 4 | 10:00:01 | main | The 10:00:00 samples land: WAL, then each series' open chunk |
| 5 | 10:00 → 12:00 | main | 720 samples and 6 chunks per series; p4's 8 series 540 samples and 5 chunks; 21,600 samples in the block |
| 6 | 12:00:11, 13:00:11 | main | The head spans more than 3 hours: blocks 08:00 to 10:00 and 10:00 to 12:00 are written; the WAL before them is dropped |
| 7 | text | main | p4's _count, paper encoder in seconds: 78, 25 and 19 bits for its first three samples |
| 8 | 12:00 | main | Bytes a sample over the block: integer counters 2.00, float _sum 6.99, gauges 0.27, all 2.41 (paper); 2.02, 6.60, 0.40, 2.39 (Prometheus). p1's _count: 704 10 XORs, 15 11, no zeros |
| J | (side) | copy | Timestamps jittered by −1, 0, +1 ms: D = 0, −3, +3; 16 bits for two timestamps in three in Prometheus's encoding (about 11 on average); p1's _count 2.03 → 3.25 bytes |
| 9 | 10:21:30 | main | Minute 10:20 rolled up: p4 _count +120, le="1" +0, le="2.5" +120, _sum +240; p1 +1,200, +1,080, +78; 44 rows |
| 10 | same | main | The minute's mean 474 ÷ 3,720 = 0.127 s; the average of pod means 0.549 s |
| 11 | same | main | p4's gauge: 0, 4, 4, 4, 4, 4 → count 6, sum 20, min 0, max 4; minute 10:21: 4, 4, 4, 4, 0, 0 |
| 12 | 11:05:00 | main | Hour 10 from 60 minute rows: _count 223,140; p4's gauge max 4, mean in flight 0.1; equal to the hour from raw for all 32 series |
| 13 | 10:41:30 | main | p2's restart at 10:40:03: last − first −119,060; reset-adjusted 1,140; 60 requests lost |
| W | (side) | copy | Decimation: the gauge keeps 4 and 0 for minutes 10:20 and 10:21 and 0 for the hour (D4's max 4); counters give the same mean (0.065775 s) and p99 (0.462 s) as D4, but 940 instead of 1,140 across the reset |
| G | (side) 10:30:02 → 10:30:05 | copy | p2 freezes for 3 s, requests counted at completion: 58 of the 60 started in the pause take over 0.1 s, 10 of them over 2.5 s, plus the 2 held in flight (3.075 s and 3.025 s). The gauge reads 2 at 10:30:00 and 10:30:10, as always: no gauge roll-up can show the pause. The histogram's minute 10:30 has 172 requests over 0.1 s instead of 120, and 12 in +Inf |
| 14 | text | main | Minute 10:20's p99 from summed buckets: 1 + (3,682.8 − 3,600) ÷ 120 × 1.5 = 2.035 s, above the slowest real latency; true 2.0 s; pods' p99s averaged 0.8 s |
| 15 | text | main | Hour 10's p99: estimate 0.462 s, true 0.4 s, D3's average of 60 minute estimates 0.503 s; the incident survives as 180 in (1, 2.5] and max 4 |
| 16 | text | main | A DDSketch at 1%: 1.994 s (0.32% off) and 0.402 s (0.62% off); 3 buckets per pod-minute (1 for p4, 2 in minute 10:21), 4 merged |
| B | (side) | copy | Bounds end at 1 s: the estimate is capped at 1.0 s, and "p99 above 1 s" can never fire |
| 17 | text | main | Unique counts don't add; HyperLogLog sketches merge (12 KB, 0.81%) |
| 18 | text | main | Values a day: raw 276,480; 1-minute 63,360; 1-hour 1,056 |
| 19 | text | main | At retention and 2 B a value: 8.3 MB, 11.4 MB, 1.5 MB; measured on our regular traffic: raw 2.41 B, roll-up values 0.36 B |
| 20 | text | main | 1.5× for a gauge, 6× for a counter, 60× only from 1 s with one value; our mix 4.4× |
| X | (side) | copy | Six designs: D0 13.8 to 27.6 MB a day; D1 grows forever; D2 loses spikes and in-window resets; D3 wrong p99 (0.503 s); D4 and D5 answer everything |
| 21 | 12:10:00 | main | 2 h at 10 s: raw, 721 points; 7 d at 5 min: 1-minute tier, 2,017 points (newest point raw); 1 y at 1 d: 1-hour tier, 366 points |
| 22 | same | main | 10:00:07 to 12:00:07 at 60 s: 121 points at :07, matching no cached answer; aligned, they match |
| 23 | 10:40 → 10:41 | main | rate() per series 18.8 then 19.0 ([1m]), 19.79 ([5m]); rate() of the summed counter 5,100.8; sum(rate()) 60.8 |
| 24 | 11:30:00 | main | p4 leaves discovery: staleness markers; absent from every evaluation from 11:30:00 |
| S | (side) | copy | A push client without markers: p4 returned by every evaluation until just before 11:34:50 |
| 25 | (side) 11:40:10 | copy | A missed scrape: rate(…[20s]) returns nothing at 11:40:25 (the first evaluation after the 11:40:20 sample lands); rate(…[1m]) has 5 samples |
| 26 | 11:12:05 | main | p3's 42 buffered scrapes land oldest first: late but in order, all 336 samples accepted |
| 27 | 11:12:30 | main | 7 dirty minutes re-rolled: 11:04 gains 200, 11:05 to 11:10 gain 1,200 each for p3 |
| 28 | 12:05:00 | main | Hour 11: 219,580 (20 of p4's requests never scraped); ignoring late data with no gate: 212,180, 7,400 short |
| O | (side) 11:12:05 | copy | Newest first: 41 rejected per series with a window of 0, 36 with 60 s, 0 with 10 minutes |
| V | (side) 12:40:05 | copy | The partition lasts until 12:40:00. The head's newest time is 12:40:00, so its appendable minimum is 11:40:00: p3's samples 11:05:00 to 11:39:50 (210 per series) are too old and go to the backfill path, and 11:40:00 to 12:39:50 are accepted in order. Hour 11 waits from 12:05 (440 rows missing). 12:40:30: 60 minutes (11:39 to 12:38) rolled again from the new samples. 13:00:00: the hourly backfill job imports the 1,680 old samples as a block for their past window. 13:00:30: minutes 11:04 to 11:39 re-rolled. 13:01:00: hour 11 rolled once, 219,580. 13:05:00: hour 12 rolled. The day's results-cache generation is bumped with each re-roll, at 12:40:30 and at 13:00:30. For those 20 minutes p3's minute 11:39 row holds 42,200: its one gap, 11:04:50 → 11:40:00, spans the whole outage, until the backfill fills it in and it drops back to 1,200 (a rejected counter sample loses no requests; its increase lands in the next accepted gap). (The metrics loop's gateway routes anything older than 10 minutes by the wall clock to the backfill path before the head is asked) |
| 29 | text | main | Roll up early and re-roll, or late and rarely (Thanos: after 40 hours): freshness against re-roll work |
| 30 | 10:20 → 10:23:30 | main | "p99 above 1 s": raw fires 10:20:30 (clears 10:22:15); the 1-minute tier 10:21:30 (clears 10:23:30); the 1-hour tier never |
| 31 | text | main | Hour 10's mean 0.0658 s against 0.0642 s in a normal hour: the incident is 0.08% of the hour |
| R | (side) 11:21:30 | copy | A retried roll-up: with add, 22 rows doubled (6,120 requests, 102 a second, for minute 11:20); with set, all 44 right |
| K | (side) 11:40 → 12:10 | copy | The worker is down: 30 minute windows missing, queries fall back to raw, hour 11 waits and is rolled at 12:10:00; run anyway it would miss 75,600 requests |
| I | (side) 11:50:00 → 11:50:30 | copy | The ingester crashes: replaying the log's 31,688 samples and markers (and 32 series entries) rebuilds all 32 series exactly; the scrapes of 11:50:00, :10 and :20 are resent and land at 11:50:30.5, in order, before the 11:50:30 scrape. They land just after the 11:50:30 roll-up run, so minute 11:49 is re-rolled at 11:51:30. Nothing acknowledged is lost |
| 32 | 10:20:10 → a year | main | One sample through every layer: client, agent, gateway, head, rule (fires 10:20:30), minute row (10:21:30), hour row (11:05), block (13:00:11), object storage, a year-long query |
The cheat card
| Topic | Remember |
|---|---|
| A series | Name + full label set; each costs about 8.3 KB of head memory (planning) and about 200 B of index |
| A sample | 16 B raw; about 1 to 2 B compressed (1.37 B was the paper's workload; our mix 2.41 B; float sums about 7 B) |
| Compression | Timestamps: delta of deltas, 1 bit when scrapes are aligned. Values: XOR with the previous, 0 / 10 / 11 |
| Head and blocks | WAL, then an open chunk of up to 120 samples; 2-hour blocks, written once the head spans more than 3 hours; retention deletes whole blocks |
| Roll-up aggregates | Counter: its reset-adjusted increase. Gauge: count, sum, min, max. Never an average without its count, never a percentile |
| Percentiles | Sum bucket counts across pods and windows, then histogram_quantile; error up to the bucket's width; top bound above the worst case; or a log-bucketed sketch within α |
| Tier size | Values a window × windows a day × days kept × bytes a value; 1.5× (gauge) and 6× (counter) from 10 s to 1 minute, not 60× |
| Tier choice | The coarsest tier no coarser than the step, covering the range, with its rows written; aggregate inside each step |
| Rates | rate() per series first, then sum; windows several scrapes long |
| Staleness and gaps | Stale is absent; no data is not zero; pair threshold rules with absent() |
| Late data | In order for its series: append. Out of order: a window (Prometheus 0, Managed Prometheus 60 s, up to 600 s). Older: backfill. Then re-roll every dirty window and bump the cache generation |
| Alerts | On a tier: window end + grace + evaluation interval. On raw: time for the bad share to cross the threshold + scrape + evaluation. Alert on raw |
| Pipeline | Set by key, never add; coarse tiers wait for every non-stale input; alarm on the oldest missing roll-up |
| AWS | CloudWatch: 3 h / 15 d / 63 d / 455 d, $0.30 a metric; Managed Prometheus: raw 150 d (up to 3 y), 60 s out-of-order window, 50M active series; Timestream for LiveAnalytics closed to new customers |
Failure checklist
- Is every label bounded, and is there an active-series limit per tenant where the series live?
- Do you size memory from series count, not from compressed bytes?
- Are scrapes aligned, so timestamps compress to a bit?
- Does every roll-up keep increases for counters and
count,sum,min,maxfor gauges, and nothing finished? - Are counter increases reset-adjusted, in roll-ups and in queries?
- Is every roll-up row set by (series, tier, window), never added to?
- Are percentiles computed only from merged buckets or sketches, never averaged?
- Do bucket bounds sit near the alert thresholds, with the top bound above the worst latency?
- Is each tier sized by the values it keeps, with a real bytes-a-value figure?
- Does every query choose its tier by step, align to the step, and report the tier it used?
- Is every late sample either appended, kept by an out-of-order window, or routed to backfill, and never re-stamped?
- Does every late sample mark its windows dirty, and does the re-roll bump the results cache's generation?
- Do coarse tiers wait for all their non-stale inputs, and do you alarm on the oldest missing roll-up?
- Do alerts read raw data, with an
absent()rule for when the data stops?
Words we use
| Word | Meaning here |
|---|---|
| Series | One metric name plus one full set of label values |
| Sample | One (timestamp, value) pair of a series |
| Counter, gauge, histogram | A running total; a reading at an instant; cumulative bucket counters plus _sum and _count |
| Summary | Quantiles computed in the client and exported finished; can't be merged |
| Head | The in-memory part of the store: the WAL, open chunks and the label index |
| Chunk, block | Up to 120 compressed samples of one series; two hours of every series' chunks plus an index, written once |
| Delta of deltas, XOR | How a timestamp's gap changed; the bits in which a value differs from the previous one |
| Roll-up, tier | Aggregates for a coarser window; all roll-ups at one resolution, with their retention |
| Increase, reset | A counter's growth over a window; a restart from 0 |
| Decimation | Keeping only the last sample of each window |
| Staleness marker | A special value written when a series ends, so queries drop it at once |
| Out-of-order window | How far behind the store's newest time a sample older than its series' newest may be and still be kept |
| Backfill | The slow path that writes data too old for the head into past blocks |
| Dirty window, re-roll | A window whose roll-up ran before all its data arrived; computing it again |
| Sketch | A small, mergeable summary of a distribution (DDSketch, t-digest) or of distinct items (HyperLogLog) |
Think-first drills
Drill 1. Five pods: four take 100 requests a second at a mean latency of 0.05 s, and a canary takes 10 a second at 1.0 s. What does the average of the pods' means say, and what is the true mean?
Drill 2. A gauge and a counter are scraped every 15 s, kept raw for 30 days, and rolled up to 5 minutes (the gauge's four aggregates, the counter's increase) kept for a year. How much does the 5-minute tier shrink one day of each, and is it bigger or smaller than raw over its retention?
Drill 3. A service fails every request for 2 minutes. An alert fires when the error ratio is above 5%. How soon can it fire on raw data with rate(…[5m]) and a 15-second evaluation, on a 1-minute tier rolled up with a 30-second grace, and on a 1-hour tier?
Interview questions
| Question | Model answer |
|---|---|
Why is a series, not a sample, the unit of cost? What happens when someone adds user_id as a label, and where do you stop it? | A sample adds a byte or two to an open chunk; a series needs an ID, label strings, posting-list entries, an open chunk and bookkeeping: kilobytes of memory (about 8.3 KB is a planning figure), whatever its sample rate. user_id multiplies series by the number of users: in our example, 32 series become 280,000 for the same traffic, 2.3 GB of head memory, and $6,000 a month as CloudWatch metrics. Stop it at the gateway (drop unbounded labels, cap counts and lengths), enforce an active-series limit per tenant where the series live (new series over the limit are dropped and counted), and alert when a tenant's series double. Estimate distinct series with HyperLogLog. |
| Walk me through how a time-series store gets a sample to about a byte. When does it do much worse? | Samples of one series sit together in a chunk. Timestamps: store the delta of deltas; with aligned scrapes it is 0, one bit. Values: XOR with the previous value; unchanged costs one bit, a small change a few meaningful bits inside the previous window (10), otherwise 11 plus 5 bits of leading zeros, 6 bits of length and the bits. It does worse when values change every sample (a busy counter, about 2 B), when a float accumulates fractions (a _sum, about 7 B), when timestamps jitter by milliseconds (16 bits instead of 1), and on each chunk's first sample. 1.37 B was the paper's workload. |
| What do you keep in a 1-hour roll-up so you can answer the mean, the maximum and the p99 of any month a year later? | Only values that merge. For each counter, including every histogram bucket, _sum and _count: its reset-adjusted increase. For each gauge: count, sum, min, max. The month's mean is Σ_sum ÷ Σ_count; its maximum is the max of the maxes; its p99 is histogram_quantile over the month's summed buckets, an estimate within the bucket's width (or within α with a log-bucketed sketch). Never a stored average or p99: averaging per-minute p99s turned a 0.4 s hour into 0.503 s. Rows are keyed by (series, tier, window) and set, so re-runs are safe. |
| A sample for 11:05 arrives at 11:12, after that minute was rolled up. What are your options, and what does each cost? | First, the store. If it is newer than its series' newest sample, it is in order and is simply appended, however late. If it is older, an out-of-order window keeps it (Prometheus's default is 0; Managed Prometheus 60 s up to 600 s; 10 minutes in the metrics loop), at a cost in memory; beyond the window it goes to a backfill path that writes past blocks, merged by vertical compaction. Then the roll-ups: mark every window it touches dirty and re-roll it, and the hour above it, with set, and bump the results cache's generation; or roll up late (Thanos waits 40 hours) and accept stale tiers; or ignore it and publish wrong history (with no completeness gate, hour 11 was 7,400 requests short). Never re-stamp it with the arrival time. |
| Why not alert on the hourly tier? How fast can each tier detect a 90-second incident? | On a tier, an alert can't fire before the window ends, its roll-up runs and the rule evaluates: our 1-minute tier fired at 10:21:30, as the 90-second incident ended. The hour dilutes it: 180 slow requests are 0.08% of 223,200, so the hour's p99 estimate is 0.462 s and never crosses 1 s. On raw data, the rule needs more than 1% of its 1-minute window slow, 38 requests, which took about 19 s, and fired at 10:20:30. Alert on raw data or recording rules; build long SLO windows from recorded short sums; use tiers for history. |
| CloudWatch, Managed Service for Prometheus or your own cluster for 100 million series: what does each keep, for how long, and what does each charge for? | CloudWatch rolls up on a fixed schedule (sub-minute for 3 hours, 1-minute for 15 days, 5-minute for 63 days, 1-hour for 455 days), accepts time stamps up to two weeks old, computes percentiles only from raw values (or Values/Counts arrays), and charges per metric a month ($0.30 falling to $0.02 at volume), so each dimension combination is a line item: at 100 million series, not a fit. Managed Prometheus keeps raw samples 150 days by default (up to 3 years), describes no roll-up tiers, has a 60-to-600-second out-of-order window and 50 million active series per workspace by default, and charges per sample ingested, per GB stored and per sample queried: about $425,000 a month in ingestion alone at 10 million samples a second. Your own Thanos or Mimir on EKS with blocks in S3 has tiers and retention you choose, and cost you size: about $76,000 a month in the metrics loop, plus the team to run it. |
Where to go next
- Write-Ahead Log, fsync & Group Commit: the head's log, and what a crash can lose (Parts 3 and 11 here).
- LSM-Trees & Compaction: time-window compaction, write stalls and corruption (Part 5 and Part 10 there), and its drill The Time-Series Database That Froze on Flush.
- Event Time, Watermarks & Checkpoints: streaming windows, the hour from closed minutes, allowed lateness and the acceptance horizon (Parts 3 and 4 there).
- Sharding, Hot Keys & Rebalancing: hashing series to ingesters, and mergeable pre-aggregation (Part 5 there).
- Idempotency & Effectively-Once Processing: absolute writes and why set beats add (Part 5 there).
- Caching & Invalidation: the results cache in general.
- Queues & Delivery Semantics: the ingest stream and consumer lag.
- Rate Limiting Algorithms: the samples-a-second limit per tenant.
- Replication, Quorums & Read-Your-Writes and Multi-Region Failover: the head's second copy, and an alerting path that survives a Region.
- Background: Write-Ahead Log & LSM-Trees, Trie & Inverted Index (posting lists), Bloom Filters & Counting Filters (sketches).
- Loops that rely on this page: metrics and alerting (the opener, R1.4, steps 1.0 to 1.2 and 1.4, R1.7, R1.9, steps 2.1 to 2.5, R2.6, R2.8, R2.9, R2.11 and steps 3.1 and 3.3 to 3.6), ad-click aggregation (steps 1.2, 1.5, 2.1 and 2.3, R1.8, R1.9, R2.5, R2.6 and step 3.4), the URL shortener (steps 2.2 and 3.5), the gaming leaderboard (step 3.2), S3-like storage (step 3.5), the job scheduler (step 3.5), payment processing (step 3.1), mobile stock trading (R1.4) and the mobile paging library (step 3.6 and R3.9).