System Design Interview

Course Content

System Design Interview

31 sections · 71 lessons

Metrics Monitoring: scope, scale and collection


"Design a metrics monitoring and alerting system." The system that tells you every other system is broken, which makes its own failure modes uncomfortable to think about.

This lesson covers scoping, the arithmetic that decides the architecture, and how data gets from a host into the pipeline. The second lesson goes inside the storage engine, the retention tiers, and the alerting layer on top.

Scope to one pillar, and say whyIn scope: metrics• Numeric time series with labels• Fixed shape, so heavy compression• Queried by dashboards and rulesOut of scope: logs, traces• Text search is a different index• Traces need per-request sampling• Each is its own hour-long design
Scoping to metrics is not dodging: the three pillars need different storage engines, and an hour fits one.

The three pillars, and why you should scope to one

Observability is conventionally split into three data types.

Metrics are numbers sampled over time: CPU at 71%, requests per second at 4,300. Small, regular, aggregatable, cheap to store.

Logs are timestamped text records of individual events. Large, irregular, searched rather than aggregated.

Traces follow one request across services, recording how long each hop took. Structured, sampled, used for latency debugging.

They have genuinely different storage engines. A design that tries to serve all three in one store serves none well. Say this, propose scoping to metrics, and offer to discuss the other two in the wrap-up. Interviewers generally accept, because the metrics half is where the interesting design is.

The questions that shape everything after

  1. Metrics only, or logs and traces too? As above — narrow it deliberately.
  2. How many hosts, and how many metrics per host? This is the scale question, and the arithmetic below shows it dominates every other decision.
  3. At what interval are metrics sampled? Every 10 seconds and every 1 second differ by 10× in every downstream number.
  4. How long is data retained, and at what resolution? "Two years" and "two years at one-second resolution" are wildly different systems.
  5. What are the query patterns? Dashboards querying the last hour, alert rules evaluating every 30 seconds, and analysts running ad-hoc year-long queries stress different parts.
  6. What is the alerting latency target? Detecting an outage within 30 seconds and within 5 minutes lead to different collection designs.

The assumptions this section uses

QuestionAssumption
ScopeMetrics only; logs and traces named as separate systems
Fleet100,000 hosts, 200 metrics each
IntervalEvery 10 seconds
Retention7 days at full resolution, 30 days at 1 minute, 2 years at 1 hour
QueriesDashboards and alert rules dominate; ad-hoc is rare
Alert latencyUnder one minute from threshold breach to notification

Requirements and scale

The write rate here is larger than anything else in Part IV (Sections 21–27), and computing it early is what justifies every subsequent choice.

The largest write rate in the subject100 M machinesPull willnot scaleabout 100 each10 B seriesevery 10 s1 B writes/s1 B per secondNo generaldatabaseabout 2 PB a dayDownsample or dieNumberConsequenceMachinesMetrics eachIntervalWrite rateRaw storage
A billion writes a second rules out every general-purpose database and forces a purpose-built time-series engine.

Functional requirements

  • Collect numeric metrics from a large fleet of hosts and services.
  • Store them with the timestamp and a set of labels identifying what they describe.
  • Serve time-range queries with aggregation: average, percentile, rate, sum by label.
  • Evaluate alert rules on a schedule and notify when they fire.
  • Power dashboards at interactive speed.

Non-functional requirements

  • Ingest must not drop data under a traffic spike, because a spike is exactly when the data matters.
  • Query latency under a second for a dashboard panel over the last hour.
  • Retention as tabled above, with cost controlled by resolution, not by deletion.
  • Available independently of the systems being monitored.

The arithmetic that decides the design

Active time series. 100,000 hosts × 200 metrics = 20 million distinct series. A series is one metric with one specific set of label values — cpu_usage on host web-041 in region ap-south-1 is one series.

Write rate. 20 million series ÷ 10-second interval = 2 million data points per second.

Hold that number next to what a relational database does. A well-tuned Postgres instance handles roughly 10,000–50,000 row inserts per second with indexes. Two million per second is 40 to 200 machines' worth of inserts, purely to write, before anyone reads anything. This is the moment a general-purpose database is eliminated, and it is worth doing in front of the interviewer rather than asserting.

Raw storage. A naive row of (timestamp 8 bytes, value 8 bytes, series identifier 8 bytes) = 24 bytes. 2 million/s × 24 bytes = 48 MB/s, or 4.1 TB per day, or 29 TB for one week at full resolution. Add labels stored per row and it doubles.

Compressed storage. Time-series compression, covered in Time-series storage, gets a data point down to roughly 2 bytes on typical data. 2 million/s × 2 bytes = 4 MB/s, 350 GB/day, 2.4 TB per week. A 12× reduction, which is the difference between a rack and a shelf.

Query load. 500 dashboards refreshing every 30 seconds with 10 panels each = 167 queries per second. Alert rules: 5,000 rules evaluated every 30 seconds = 167 per second. Both modest. The read side is not the problem; the write side is.

What the numbers rule out

  • A relational database: 40–200 machines to absorb the inserts, and its B-tree index does random writes, which is the wrong access pattern entirely.
  • Storing every point forever at full resolution: 29 TB/week uncompressed is 1.5 PB over a year, for data nobody queries at that resolution after a week.
  • One node: 48 MB/s is fine for one disk, but 20 million series will not fit one node's memory index, and the failure domain is unacceptable.

Push versus pull collection

With the scale fixed, the first box on the diagram is collection. How does a data point get from a host into the pipeline? Two models, both used in production by serious systems, with a genuine trade-off and no universal answer.

Push

An agent on each host samples locally and sends batches to a collection endpoint. The monitoring system is passive; it accepts what arrives.

Advantages. Works through firewalls and network address translation, because the connection is outbound from the host. Works for short-lived things — a batch job that runs for 40 seconds can push its results before exiting. No service discovery needed, since hosts announce themselves by sending.

Disadvantages. A silent host is ambiguous: is it healthy and idle, dead, or misconfigured? You cannot tell from absence. A misbehaving agent can flood the collector, so the collector needs its own rate limiting per source. And the collector is a write endpoint exposed to everything, so it needs authentication at scale.

Pull

The monitoring server holds a list of targets and scrapes an HTTP endpoint on each one on a schedule. The host exposes current values; the server decides when to read.

Advantages. A failed scrape is an unambiguous signal — the server knows it tried and failed, so "target down" is a first-class metric you get for free. The server controls the rate, so it cannot be flooded. You can scrape a target manually from a laptop to debug it, which is genuinely useful at 3 a.m.

Disadvantages. Requires service discovery: something must maintain the list of what to scrape, usually by integrating with the orchestrator or cloud provider inventory. Requires inbound network reachability to every target, which is awkward across network address translation or strict firewalls. Short-lived jobs may finish before any scrape happens, which needs a separate push gateway — an exception that partially reintroduces the push model.

SOURCESApp serversDatabasesHosts, networkPull collectorscrapes /metricsPush gatewayagents send inIngest queueTime-series DBhot, 15 daysRule evaluatorruns every 30 sDashboardsAlert managerdedupe, group, routeOn-callscrapequeryfirespull knows when a target disappears;push does not, which is why pull is thedefaultgrouping and dedupe live here, not in therules — one incident should page once
Pull collection gets you liveness for free: a target that stops answering is itself the signal.

The recommendation, and the honest caveat

Pull for a fleet you control, because the dead-host signal and the rate control are worth the service-discovery work, and both problems get harder as the fleet grows. Push for anything ephemeral, external, or behind a network boundary — serverless functions, batch jobs, customer-hosted agents.

Almost every real deployment runs both. Pull as the backbone, a push gateway at the edges. Presenting this as a hybrid rather than a doctrine is the answer that scores, because the interviewer has usually operated one of these and knows the exceptions exist.