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.
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.
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.
How a rule is evaluated
An alert rule is a query plus a condition plus a duration:
error_rate{service="checkout"} > 0.05 for 5mA 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_downin rackr14". - 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_downis firing, suppress everyhost_downin 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.