Mid to Senior Engineer

System Design Interview Prep

A structured path from the interview framework through core concepts, key technologies and patterns to eighteen full problem breakdowns, each with diagrams and weak, solid and excellent answers to every deep dive.

Chapter 29 of 36Problem breakdowns · Design a Metrics and Monitoring System

Design a Metrics and Monitoring System

A monitoring system collects numbers from thousands of machines and services, stores them cheaply, lets engineers query and graph them, and pages someone when something is wrong. It is a useful interview problem because it is write-heavy and read-light in a way most systems are not, it relies on a specialised storage model, and it contains a famous scaling trap, called cardinality.

The chapter follows the usual shape: understand the problem, set up the interface, build the high-level design, then go deep on the questions interviewers use to separate levels.

1. Understanding the problem

Services and machines emit metrics: numeric measurements over time, such as CPU usage, request count, queue depth and error rate. The monitoring system gathers them, keeps a history, answers queries about them, and raises alerts when conditions are met.

Functional requirements

Core:

  1. Collect metrics from many sources, each sample being a name, a set of labels, a value and a timestamp.
  2. Store them and allow queries over time ranges with aggregation, such as "the 95th percentile latency of the checkout service over the last six hours, per region".
  3. Evaluate alert rules and send notifications when they fire.

Confirm in or out: logs and traces (usually separate systems), dashboards, long-term retention, anomaly detection, and multi-tenant isolation. A sensible opening: "I will design metrics collection, storage, querying and alerting, and treat logs and traces as separate systems."

Non-functional requirements

  • High write throughput. Millions of samples arrive every second, steadily, forever.
  • Low query latency for dashboards, within a second or two for typical queries.
  • Reliability of alerting. Monitoring must keep working when the systems it watches fail, so it is among the most highly available components.
  • Cost efficiency. The volume is enormous, so storage must be compact and old data must be reduced.
  • Freshness. Data should be queryable within seconds, and alerts should fire within a minute or so.

Estimation

Assume 1 million machines or containers, each emitting 100 metrics every 10 seconds.

QuantityCalculationResult
Sample rate samples per second
Raw size16 bytes per sample uncompressed160 MB per second
Per dayabout 14 TB per day uncompressed
Compressedspecialised encodings often reach 10x or moreabout 1 to 2 TB per day

What the numbers say. Ten million writes a second sustained is beyond a general-purpose database, so the design needs a write-optimised, append-only store and a buffer that absorbs bursts. Reads are few by comparison, but they scan many points, so storage layout matters. And because 14 TB a day cannot be kept forever at full resolution, downsampling is essential.

2. The set up

Core entities

  • Metric series: a name plus a unique combination of labels, such as http_requests_total{service="checkout", region="eu", status="500"}.
  • Sample: a timestamp and a value in a series.
  • Alert rule: a query, a condition, a duration and a destination.

Interfaces

Two models for collection:

  • Pull: the monitoring system periodically scrapes an endpoint on each target. It makes it easy to know whether a target is alive (the scrape fails), and the system controls the rate. It needs service discovery to find targets.
  • Push: each target sends its metrics to a collector. It suits short-lived jobs and environments where the monitor cannot reach targets. A push gateway accepts metrics from batch jobs.

Many systems support both. State which you would choose and why.

A query interface accepts a series selector, a time range, a resolution and aggregation functions:

query: sum by (region) (rate(http_requests_total{service="checkout"}[5m]))
range: last 6 hours, step 1 minute

3. High-level design

The pipeline has five stages:

  1. Agents on each host or in each service emit metrics. Collectors receive them, validate and tag them.
  2. An ingest log, a durable queue, buffers the samples. It absorbs bursts, decouples collection from storage, and lets several consumers read the same stream.
  3. Writers batch the samples, compress them and write them to the time-series store, sharded by series.
  4. A query service reads the store, with a result cache, for dashboards and API calls.
  5. An alert evaluator runs the rules against fresh data and sends notifications.
<!--fig:hld-->
fire query Agents on each host / service Push gateway batch jobs Collectors validate, tag Ingest log Writers batch + compress Alert evaluator rules on fresh data Time-series store sharded by series Query service + result cache Notifications page, chat, email Figure 1. Collect, buffer, store, query and alert. The log between collection and storage absorbs bursts and lets alerting and storage consume independently.

Notice that alerting reads from the stream and from recent data, not through the slow path of the dashboards. That keeps alerts fast and keeps them independent of query load.

4. Potential deep dives

Deep dive 1: How do you store ten million samples a second?

The challenge. A general-purpose database cannot absorb this write rate, and the access pattern is unusual: data arrives in time order, is almost never updated and is read by time range for one series at a time.

Weak: a relational table with a row per sample. Ten million inserts a second, each with an index update, overwhelms a relational database. The row overhead dwarfs the 16 bytes of useful data, and the table grows without bound.

Solid: a time-series layout, partitioned by series and by time. Group samples by series, store each series' points together in time order, and partition data into time blocks, such as two-hour chunks. Writes append to an in-memory buffer backed by a write-ahead log, then flush to an immutable block. Queries touch only the blocks in their time range and read contiguous data for each series.

Excellent: columnar-style compression, in-memory head and sharding. Exploit the regularity of the data. Timestamps arrive at near-regular intervals, so store the delta of deltas, which is mostly zero and compresses to a few bits. Consecutive values in a series are often similar, so store them as the XOR of adjacent values or similar delta encodings, which often compress a sample from 16 bytes to a couple of bytes. Keep the most recent data in memory for fast writes and queries, and flush it to compressed blocks on disk periodically. Shard by series identifier, using consistent hashing, so that each series lives on one node and its points stay contiguous, and replicate each shard for durability. Explain that the append-only, immutable-block design avoids random writes and makes compaction and retention simple: deleting old data is dropping whole blocks.

Deep dive 2: How do you handle retention and old data?

The challenge. Fourteen terabytes a day uncompressed becomes petabytes quickly, and few people need ten-second resolution for data from last year.

Weak: keep everything at full resolution forever. Storage cost grows without limit, and queries over long ranges become slow because they read millions of points.

Solid: delete data older than a retention period. Simple, and it loses history that is occasionally valuable for capacity planning and year-over-year comparison.

Excellent: downsample into tiers, then delete. Keep raw samples for a short period, then roll them up into coarser resolutions that retain the useful statistics: minimum, maximum, average, sum and count per interval.

<!--fig:tiers-->
roll up roll up Raw 10-second sampleskept 15 days 1-minute rollups min, max, avg, countkept 90 days 1-hour rollups kept 2 years Queries pick the coarsest tier that still gives enough points for the requested time range. Figure 2. Downsampling keeps recent data precise and old data cheap: roll up before you delete.

Queries choose the coarsest tier that still gives enough points for the requested range, so a graph over a year reads hourly points and a graph over an hour reads raw ones. Roll up before deleting, and keep the rollups needed to compute accurate aggregates. Note a subtlety: an average of averages is wrong unless weighted by counts, and percentiles cannot be recomputed from rollups of averages, so store histograms or sketches for latency metrics if you need long-term percentiles.

Deep dive 3: The cardinality problem

The challenge. The number of unique series, called cardinality, determines memory, index size and cost. It multiplies with every label.

The trap. A metric with the labels service (100 values), region (10), status (5) has series, which is fine. Add a label such as user_id with a million values and the same metric becomes five billion series. Each series has an index entry and in-memory state, so the system runs out of memory and falls over. This is the single most common way monitoring systems are taken down, usually by a well-meaning engineer.

Weak: allow any label. One mistake can bring down the system for everyone.

Solid: guidelines and limits. Document that labels must have a bounded set of values, and enforce a maximum number of series per metric and per tenant. Reject or drop samples above the limit, and alert the owner.

Excellent: guard rails, visibility and the right tool. Enforce limits at the collector, track the top metrics by series count on a dashboard, and provide per-team quotas. Educate that high-cardinality identifiers such as user, request or session identifiers belong in logs or traces, which are built for them, not in metrics. When detail is needed, use sampling or aggregate before sending. Index structures should be built for lookups by label, such as an inverted index from label values to series.

Deep dive 4: How does alerting stay reliable?

The challenge. An alerting system that fails silently is worse than none, because people assume that no alert means no problem.

Weak: evaluate alerts by running queries against the main store. If the store is overloaded or down, alerts stop, exactly when you need them.

Solid: a dedicated evaluator with its own data path. Run the alert evaluator on fresh data from the ingest stream or a recent in-memory window, separate from dashboard queries. Rules have a for-duration, such as "error rate above 5 percent for 5 minutes", which prevents a single blip from paging.

Excellent: highly available evaluation with deduplication and a watchdog. Run the evaluators redundantly in several zones, and deduplicate notifications so that redundancy does not mean double pages. Group related alerts, so one outage that triggers a hundred alerts produces one notification. Add silencing for planned maintenance and inhibition, so that a "datacenter down" alert suppresses the thousand alerts it implies. Include a dead man's switch: an alert that should always be firing and that an external system watches, so if the monitoring system itself dies, the absence of its heartbeat pages someone. Route notifications by severity and team, with escalation if nobody acknowledges.

Say plainly that good alerts are symptom-based, for example high error rate or latency that users feel, and that alerting on every possible cause creates noise people learn to ignore.

Deep dive 5: Query performance

The challenge. A dashboard may run dozens of queries every few seconds across many users.

Weak: every panel queries the store each time. The store is overloaded by repeated identical queries.

Solid: cache results and limit query cost. Cache query results with a short TTL, and cache partial results by time block so that only the newest block is recomputed. Limit the maximum range, series count and duration of a query.

Excellent: pre-aggregation and query planning. For expensive, commonly used queries, compute recording rules that precompute the result on a schedule and store it as a new, small series, so dashboards read a cheap series. Push computation close to the data, aggregating on each shard before combining, which reduces data moved across the network. Choose the downsampling tier by range automatically, and show the resolution used so users are not misled.

Deep dive 6: The monitoring system's own failure

Monitoring must outlive what it monitors. Run it in separate infrastructure and, ideally, a separate failure domain, so a regional outage does not blind you. Buffer data at the agents during a collector outage and send it later, with a bounded buffer so that agents do not exhaust memory. Make writes idempotent by deduplicating on series and timestamp, so retries are harmless. Replicate storage, and monitor the monitor with a simpler independent system.

5. What is expected at each level

Mid-level. You describe agents, a storage layer and a query interface, and recognise that a relational database is a poor fit for the write rate. You mention alerting.

Senior. You size the system from the estimates, use a queue to buffer, explain time-series storage, partitioning by series and time, compression ideas and downsampling. You identify cardinality as a risk and design alerting to be reliable.

Staff. You discuss multi-tenant isolation and quotas, cost control, how recording rules and rollups interact with accuracy, the reliability of the monitoring system itself, an SLO-based alerting philosophy, and how teams are encouraged to instrument well. You decide where metrics end and logs and traces begin.

6. Interview questions and model answers

Q: Pull or push? Pull gives control over rate and a built-in liveness signal, but needs discovery and network access to targets. Push suits short-lived jobs and restricted networks. I would support both: pull for long-running services and a push gateway for batch jobs.

Q: Why not use a relational database? Ten million writes a second of append-only, time-ordered data are a poor fit for row-based storage and per-row indexes. A time-series store partitions by series and time, compresses regularly spaced values heavily, and drops old data by deleting whole blocks.

Q: What is cardinality and why does it matter? The number of unique series, which multiplies across labels. A label with unbounded values, such as a user identifier, creates billions of series and exhausts memory. I enforce limits and keep such identifiers in logs and traces.

Q: How do you keep years of data affordable? Roll up raw samples into coarser tiers, keeping minimum, maximum, sum and count, and delete the raw data after its retention. Queries choose the coarsest tier that gives enough points.

Q: How do you stop alert storms? Group related alerts, deduplicate across redundant evaluators, inhibit alerts implied by a larger failure, and require a condition to hold for a duration before firing. Alert on symptoms, not every cause.

Q: What happens if monitoring fails? A dead man's switch alert watched by an external system, redundant evaluators in several zones, agent-side buffering during outages and a separate failure domain for the monitoring stack.

7. Common mistakes

  • Storing samples as rows in a general-purpose database.
  • Keeping all data at full resolution forever.
  • Allowing unbounded label values, which causes a cardinality explosion.
  • Running alerts through the same path as dashboards.
  • No grouping or inhibition, so one outage produces a flood of pages.
  • Averaging averages when rolling up.
  • Monitoring that fails together with what it monitors.
Header Logo