System Design Interview

Course Content

System Design Interview

31 sections · 71 lessons

Metrics Monitoring deep dive: time-series storage, downsampling and alerting


The first lesson ended with 2 million data points per second arriving from a pull-based collection layer, and a relational database ruled out. This lesson builds what replaces it: a time-series store, a retention plan that makes two years affordable, and the alerting and dashboard layer that people actually see.

A time-series database is not a relational database with a timestamp column. The differences are structural, and they are what makes 2 million writes per second affordable.

The three structural properties

Writes are append-only and nearly time-ordered. New points always have the newest timestamp. So the storage engine never has to insert into the middle of an index, which is what makes a B-tree slow under this load. Log-structured storage — the same idea as Storage engines: LSM tree versus B-tree — is the natural fit.

Reads are range scans over one series, or over many series with the same label. Never "find the row with value 71.4". So the index maps labels to series identifiers, and the data for one series is stored contiguously so a range read is one sequential scan.

Old data is never updated. Immutability makes compression, compaction, and tiering all straightforward.

Compression, computed

This is where the 12× comes from, and it is two separate tricks.

Timestamps: delta-of-delta. Points arrive every 10 seconds. Store the first timestamp in full (8 bytes), then the difference from the previous one (10), then the difference between consecutive differences — which is 0 whenever the interval is regular. A zero encodes in a single bit. A perfectly regular series costs about 1 bit per timestamp instead of 64.

Values: XOR against the previous value. CPU usage goes 71.2, 71.4, 71.3. As 64-bit floating point numbers these differ only in their low bits, so the exclusive-or of consecutive values has long runs of leading and trailing zeros. Store the count of leading zeros and the few meaningful bits in between.

Facebook's published description of its Gorilla in-memory time-series store reports an average of around 1.37 bytes per data point on real production data using this pair of techniques — worth citing, and worth flagging as their measurement on their workload rather than a universal constant.

At 2 bytes per point, our 2 million points per second is 4 MB/s. At 24 bytes it is 48 MB/s. Over two years the difference is roughly 250 TB against 3 PB.

One metric, three retention tiersraw10-second pointskept 15 daysstorage 1.0×downsample5-minute rollupmin / max / avg / countkept 6 monthsstorage 0.03×downsample1-hour rollupmin / max / avg / countkept 2 yearsstorage 0.003×the row layoutA point is (metric, label set, timestamp, value). Points for one seriessit next to each other and timestamps are delta-encoded, so a day of10-second data compresses to a couple of bytes per point.why downsample at allNobody looks at 10-second resolution from eight months ago, buteveryone wants the yearly trend. Keeping both at full resolution costs30× more for a view nobody opens.Storing min and max alongside the average is what keeps a rollup honest — an average alone hides exactly the spike you are looking for.
Each tier is 30× smaller than the one above, which is what makes two years of history affordable.

Downsampling and retention

Compression shrinks each point. Downsampling decides how many points you keep at all. Nobody looks at 10-second resolution data from eight months ago; storing it is paying full price for a resolution no query uses.

Three tiers, each a fraction of the cost10 s raw, kept 7 days1 min rollup, 30 days1 hour rollup, 1 yearPercentiles cannot re-add
Downsampling is lossy in a specific way: you can average averages, but you cannot recover a p99 from stored means.

The three tiers, costed

Take the compressed figure of 2 bytes per point and 20 million series.

Tier 1 — 7 days at 10-second resolution. 2 million points/s × 2 bytes × 604,800 s ≈ 2.4 TB. This is the incident-debugging tier: when something broke at 03:14, you need to see the spike.

Tier 2 — 30 days at 1-minute resolution. Downsampling 10-second data to 1-minute is a 6× reduction: 20 million series × 1 point/min × 43,200 min × 2 bytes ≈ 1.7 TB for 30 days. This is the "is this week worse than last week" tier.

Tier 3 — 2 years at 1-hour resolution. 20 million series × 24 points/day × 730 days × 2 bytes ≈ 700 GB. This is the capacity-planning and annual-trend tier.

Total across all three: roughly 4.8 TB, replicated three ways for 14.4 TB.

Compare against keeping everything at full resolution for two years: 2 million/s × 2 bytes × 63 million seconds ≈ 250 TB, or 750 TB replicated. Downsampling is a 50× saving, and it is the single largest cost decision in the system.

What downsampling costs you in accuracy

Aggregation is lossy, and which loss you get depends on which aggregate you keep.

Keep only the average and a one-second spike to 100% CPU inside a one-minute window becomes a 2% bump in the average — invisible. Keep the maximum and the spike survives, but every window now looks spiky.

Recommendation: store four aggregates per downsampled point — minimum, maximum, sum, and count. Sum and count reconstruct the average and let you re-aggregate correctly over longer spans; minimum and maximum preserve the outliers. Four values at 2 bytes each is 8 bytes per downsampled point instead of 2, which raises Tier 2 and Tier 3 to about 9.6 TB total — still a 26× saving and far more useful data.

Alerting and visualisation

Storage was the hard engineering. Alerting is where the system succeeds or fails operationally, and the failure modes are human as much as technical.

From a rule to a page at 3 a.m.Rule queriesthe rollupThresholdbreachedHeld forthe durationGrouped,deduplicatedRouted, orsilencedOne switch failure can breach a thousand rules at the same instant.
Grouping and silencing exist because the failure mode is a human ignoring the hundredth page, not a missed query.

How a rule is evaluated

An alert rule is a query plus a condition plus a duration:

Text
error_rate{service="checkout"} > 0.05  for 5m

A scheduler runs the query every evaluation interval — 30 seconds in our design — against the time-series database. If the condition holds, the rule enters a pending state. Only after it has held continuously for the for-duration does it become firing and generate a notification.

That for-duration is the single most important field. Without it, a 30-second blip at 3 a.m. wakes someone for a problem that has already resolved. With for: 5m, the condition must persist through ten consecutive evaluations. The cost is five minutes of detection delay, which is why the value is tuned per rule: for: 1m on a total outage, for: 15m on a slow disk filling up.

Grouping, deduplication, and silencing

A rack loses power and 40 hosts go down. Forty rules fire. Forty pages is not forty times as useful as one.

  • Grouping collapses alerts sharing labels into one notification: "40 instances of host_down in rack r14".
  • Deduplication ensures a rule that is still firing at the next evaluation does not send again — one notification per alert instance until it resolves.
  • Inhibition suppresses dependent alerts when a cause alert is already firing: if rack_power_down is firing, suppress every host_down in that rack.
  • Silencing lets an operator mute a matcher for a window during planned maintenance.

Escalation then routes by severity and time: page the on-call engineer, escalate to the secondary after 10 minutes unacknowledged, escalate to the manager after 30.

Dashboards must not query raw data

A dashboard panel showing 30 days of request rate across 5,000 hosts, computed from 10-second data, touches 5,000 × 259,200 = 1.3 billion points. That query takes tens of seconds and pins CPU on the database.

The same panel against 1-minute downsampled data touches 216 million points; against pre-aggregated per-service rollups it touches a few thousand. Route dashboard queries to the coarsest tier that satisfies the requested time range, and precompute the aggregates that dashboards actually ask for — a recording rule that continuously computes sum by (service) (rate(http_requests_total[5m])) turns a fan-out over 5,000 series into a read of one.