SysDesignPrep.com
System design interview question

Design a Metrics and Monitoring System

Ingest 10 million samples a second from 100,000 hosts, keep them queryable for a year at shrinking resolution, draw dashboards in under a second, and page the right person within a minute.

Last updated 2026-09-30. Difficulty: hard. Patterns: time-series, compression, cardinality, downsampling, alerting. Reported at Datadog, Google, Amazon, Microsoft, Uber.

Walk through a strong candidate's answer, turn by turn.

The interviewer asks, the candidate answers and draws, and you press Next. Pause to answer yourself at the key decisions, and ask the AI Mentor anything along the way.

Functional requirements

  • Collect metrics. Counters, gauges and histograms from hosts, containers and applications, each identified by a name and labels (service="api", region="eu-west-1").
  • Query and dashboards. Aggregate across labels and over time ("p99 latency by region, last 6 hours"), rendered as graphs that refresh every 30 seconds.
  • Alerting. Rules over queries ("error rate above 2 % for 5 minutes") that notify on-call through PagerDuty or Slack, grouped and deduplicated.
  • Retention. Full resolution for two weeks, coarser rollups for a year.
  • Multi-tenant. Many teams (or customers) share the platform with per-tenant limits.
  • Out of scope. Logs and traces (mention how they link), anomaly detection models, and the dashboard UI.

Non-functional requirements

  • Ingest scale (10 M samples/s · 100 M active series). 100 k hosts × ~1,000 series each, reported every 10 seconds.
  • Alert latency (page within 1–2 min of the condition). This is the requirement that drives the hardest trade-off: alerting must keep working when the systems it watches, and parts of the monitoring system itself, are failing, so the alert path has to be simpler and more available than everything else.
  • Query latency (p95 < 1 s for a dashboard panel). Recent data comes from memory; long ranges from pre-aggregated rollups, never from raw samples.
  • Write availability (no gaps during an incident). Monitoring is most needed when everything else is on fire. Ingest must absorb bursts and degrade by dropping detail, not by going dark.
  • Cost (~1.4 bytes per sample stored). At 10 M samples/s, a few bytes per sample is the difference between terabytes and petabytes a month.
  • Isolation. One team adding a high-cardinality label must not take down ingestion or queries for everyone else.

Back-of-envelope estimates

  • Samples per second: ~10 M. 100 k hosts × 1,000 series ÷ 10 s interval = 10 M samples/s.
  • Raw size without compression: ~14 TB/day. A sample is a timestamp and a float, 16 bytes: 10 M × 16 B × 86 400 ≈ 13.8 TB/day. Storing that raw is not an option.
  • Compressed size: ~1.2 TB/day. Delta-of-delta timestamps and XOR-encoded floats (Gorilla) average ~1.37 bytes per sample: 10 M × 1.37 × 86 400 ≈ 1.2 TB/day, about 12× smaller.
  • Raw retention (15 days, 3 replicas): ~54 TB. 1.2 TB × 15 days ≈ 18 TB; with replication factor 3 in the ingest path and blocks deduplicated in object storage, roughly 54 TB of hot storage at peak, much less once blocks are compacted.
  • Dashboard queries per second: ~7 k. 10 k dashboards open × 20 panels ÷ 30 s refresh ≈ 6.7 k queries/s, most for the last hour, which is served from memory.
  • Alert rule evaluations per second: ~33 k. 1 M alert rules evaluated every 30 s = 33 k/s. Each is a query over the last few minutes; rules are sharded across evaluators.

Components

  • Agents (on every host): Collect host and application metrics locally, pre-aggregate where possible, buffer to disk if the network is down, and send compressed batches every 10 seconds.
  • Dashboards: Graphs that run queries over a time range and refresh every 30 seconds. The largest source of read load.
  • Ingest gateway (auth · limits): Authenticates tenants, validates names and labels, enforces per-tenant rate and cardinality limits, and writes accepted samples to the ingest log, partitioned by series.
  • Ingest log (Kafka · key = series hash): A durable buffer between agents and storage. Absorbs bursts, lets ingesters be restarted without losing data, and lets a second consumer (rules, anomaly detection) read the same stream.
  • Ingesters (in-memory head + WAL): Each owns a shard of series. Appends samples to compressed in-memory chunks, writes a write-ahead log, serves queries for the last few hours, and every two hours cuts an immutable block.
  • Series index (inverted index: label → series ids): Maps label pairs (service="api") to the series that have them, so a query can find its few hundred series among 100 M without scanning. Built per block and in memory for the head.
  • Block storage (object storage · 2 h blocks): Immutable blocks of compressed chunks plus their index, in S3. Cheap, durable, and read by queries for anything older than a few hours.
  • Compactor (merge · downsample): Merges small blocks into larger ones, removes the duplicate copies written by replicated ingesters, and produces 5-minute and 1-hour rollups (min, max, sum, count) for long-range queries.
  • Query engine: Parses a query, picks a resolution from the time range, fans out to ingesters for recent data and to block stores for older data, merges, and aggregates. Splits long ranges into cached day-sized pieces.
  • Query cache: Caches results per query and aligned time step, so a dashboard refreshing every 30 seconds only computes the newest interval.
  • Rule evaluator (sharded by rule): Runs every alert rule on a fixed interval against recent data, keeps per-alert state (pending, firing, resolved), and sends firing alerts to the alert manager.
  • Alert manager (dedupe · group · route): Deduplicates alerts from replicated evaluators, groups related ones into a single notification, applies silences and inhibitions, and routes to the owning team’s channel.
  • PagerDuty / Slack: Where humans are notified. Outside our control and rate-limited, so the alert manager batches and retries.
  • Meta-monitoring (independent stack): A small, separate monitoring setup that watches the monitoring system itself, including a "dead man’s switch" alert that must arrive every minute; if it stops, the alerting path is broken.

User flows

  1. A sample travels from a host to storage. The write path that runs 10 million times a second: batch at the edge, validate and limit, buffer durably, compress in memory, flush immutable blocks.
    1. The agent sends a compressed batch of samples every 10 seconds. Batching and compression cut 10 M samples/s into ~10 k requests/s. If the network is down, the agent buffers to local disk for a while and replays later, so short outages leave no gaps.
    2. The gateway validates labels and enforces the tenant’s rate and series limits. A tenant over its active-series limit has new series rejected (existing ones keep flowing), and the response tells the agent which metric caused it. This is the main defence against cardinality explosions.
    3. Accepted samples are written to the ingest log, partitioned by a hash of the series’ labels. Hashing by series sends every sample of a series to the same partition and therefore the same ingester, which is what lets it compress consecutive samples together.
    4. An ingester appends each sample to its series’ in-memory chunk and to its write-ahead log. Timestamps are stored as delta-of-delta (usually zero for a regular 10 s interval, so one bit) and values as XOR against the previous value, averaging ~1.4 bytes per sample. A new series is added to the in-memory inverted index.
    5. Every two hours the ingester cuts an immutable block and uploads it to object storage. A block holds compressed chunks for its series and their index. Once uploaded, the ingester can drop that data from memory and its WAL segments.
    6. The compactor merges blocks, drops replica duplicates and writes rollups. Two-hour blocks are merged into daily and then two-week blocks, which makes long queries open far fewer files. 5-minute and 1-hour rollups store min, max, sum and count per series so any aggregation stays correct.
  2. Load a dashboard. Find the right series through the index, read recent data from memory and older data from blocks at a resolution that fits the screen, and cache by time step.
    1. A panel asks for p99 latency by region over the last 6 hours.
    2. The cache returns everything except the newest step; only that piece is computed. Results are cached per query and aligned time step. A dashboard refreshing every 30 s recomputes 30 s of data, not 6 hours.
    3. The query engine uses the inverted index to find the matching series. Intersecting posting lists for service="api" and the metric name finds a few thousand series out of 100 M, without scanning anything.
    4. Recent hours are read from ingesters; older data from blocks, at a resolution chosen from the range. A 6-hour graph 1,000 pixels wide needs at most one point per ~20 s, so raw data is fine; a 90-day graph uses 1-hour rollups. Reading raw samples for a year-long graph would touch billions of points to draw a thousand.
    5. Partial results are merged and aggregated, and replica duplicates removed. Each series is held by three ingesters, so the engine deduplicates by series and timestamp. Aggregation (sum by region, quantiles) happens after merging.
  3. Someone adds user_id as a label. The scale-breaking case in every metrics system: series count explodes, and with it memory, index size and query cost. Limits contain it to one tenant.
    1. A deploy adds a user_id label to a request counter. Each distinct label combination is a separate series. With 5 M active users, one metric becomes 5 M series per host group: more than the rest of the tenant combined.
    2. The gateway sees the tenant’s active series approaching its limit and rejects new series for that metric. Existing series keep flowing, so dashboards do not go blank. The team gets an error naming the metric and label, and an alert.
    3. Without limits, ingesters would run out of memory and the index would balloon. Every series costs memory for its head chunk and index entries (a few KB). Millions of new series in minutes can crash ingesters, which then lose their unflushed data for every tenant on them.
    4. Queries touching the metric would scan millions of series. Per-query limits on series and samples fail such queries quickly instead of tying up the query fleet.
  4. An ingester crashes. The failure path: replicated writes keep queries complete, and the WAL plus the ingest log bring the crashed shard back without gaps.
    1. One ingester dies with two hours of unflushed data in memory. Its series are also held by two other ingesters (replication factor 3), so queries still see every sample; the engine just gets results from two replicas instead of three.
    2. A replacement starts, replays its write-ahead log, then resumes consuming from its last committed offset. The WAL restores the in-memory chunks; Kafka fills the gap between the WAL and now. Offsets are committed only after samples are in the WAL, so nothing falls between them.
    3. Meta-monitoring notices the ingester was down and records the event. The monitoring system cannot reliably watch itself, so a small independent stack scrapes its health.
    4. The compactor later removes the duplicate copies in uploaded blocks. Three ingesters upload three copies of the same series; vertical compaction merges them into one, so storage does not triple in object storage.
  5. An alert fires and someone is paged. The background path that matters most: evaluate, wait out blips, deduplicate, group, route, and prove the path itself is alive.
    1. Every 30 seconds the rule evaluator runs "5xx rate above 2 % for 5 minutes" for each service. Rules are sharded across evaluators by hash, and each shard runs on two evaluators for redundancy. They query only recent data, which is in ingester memory.
    2. The condition holds; the alert goes pending, then firing after five minutes. The "for" duration filters out blips. Each alert’s state is kept per label set (service="checkout"), so different services fire independently.
    3. Firing alerts are sent to the alert manager, which deduplicates, groups and routes them. Both evaluators send the same alert; the alert manager keeps one. Twenty services failing because of one database become one grouped notification. Silences and inhibitions (do not page for symptoms when the cause already paged) apply here.
    4. The owning team is paged. Routing uses labels (team="payments"). Notifications are retried and repeated until acknowledged.
    5. A heartbeat alert that always fires proves the path works; its absence pages someone. The dead man’s switch: an alert that fires every minute and is routed to an external service that pages if it stops arriving. It catches the failure mode where alerting is silently broken.

Deep dives

Push or pull

Should the monitoring system scrape targets, or should agents push to it?

Pull (Prometheus-style scraping) gives the monitoring system control of the rate and a built-in health check: if a scrape fails, the target is down. It needs service discovery to know every target, and every target must be reachable from the scrapers.

Push lets short-lived jobs, devices behind NAT and customers’ hosts send data without being reachable, and puts the batching and buffering on the agent. It loses the free liveness signal and needs strong ingest-side limits, because anyone with a key can send anything.

  • Agents on each host scrape local targets (pull) and push batches to a central ingest (push) chosen
  • Central pull (Prometheus servers scraping everything) situational: a single cluster or a Kubernetes environment with service discovery
  • Applications push directly to the backend situational: serverless functions and short-lived jobs

The answer: An agent on every host discovers local targets, scrapes them every 10 seconds (recording up/down per target as its own metric), pre-aggregates where configured, and pushes compressed batches to the ingest gateway with disk buffering for outages. Short-lived jobs push through the agent or a push endpoint. The agent’s own heartbeat is a metric, alerting on its absence. This is the architecture of Datadog’s agent and of Prometheus agents remote-writing to a central store.

How do you tell "host is down" from "agent stopped reporting"?

From the backend alone, you cannot. Alert on absence of the agent heartbeat, then cross-check with an external signal: the cloud provider’s instance state, a ping from another host, or load balancer health checks. The alert should say which signals agree.

Why scrape every 10 seconds rather than every second?

Cost scales linearly with frequency: 1 s would be 100 M samples a second. Ten seconds resolves most incidents well. For the few metrics that need more, histograms and higher-frequency local aggregation capture short spikes without storing every second.

How do agents avoid all sending at the same second?

Each agent offsets its 10-second schedule by a random amount at startup, so 100 k hosts spread their pushes evenly across the interval. Without that jitter, the gateway would see a spike every 10 seconds and idle in between.

What happens when the gateway sheds load?

It returns a retryable error and the agent keeps the batch in its disk buffer, retrying with backoff. Data arrives late but complete. Only when the buffer fills does the agent drop the oldest data, and it reports how much it dropped as a metric.

Storing time series cheaply

Ten million samples a second. Why not just insert rows into a database?

As rows, each sample is a timestamp, a value and a series reference: 10 M inserts a second, about 14 TB a day before indexes. A general-purpose database would spend most of its effort on index maintenance for data that is written once and read in time-ordered ranges.

Time series are extremely regular. Timestamps arrive at fixed intervals, and consecutive values of a gauge or counter are close to each other. Storing differences instead of values, at the bit level, shrinks samples by an order of magnitude.

  • Purpose-built TSDB: per-series compressed chunks in memory, WAL, immutable blocks with an inverted label index, in object storage chosen
  • Wide-column store (one row per series per time bucket) situational: teams already running Cassandra or Bigtable at scale (as OpenTSDB did on HBase)
  • Relational database with a timestamp index rejected

The answer: Each ingester keeps an in-memory "head" per series: chunks of ~120 samples encoded with delta-of-delta timestamps and XOR-compressed floats (Facebook’s Gorilla encoding), plus a write-ahead log for crash recovery. Every two hours, the head is cut into an immutable block containing the chunks and an inverted index from label pairs to series ids. Blocks are uploaded to object storage; queries read them through caching store gateways. The compactor merges blocks over time. Series are sharded across ingesters by a hash of their labels, with replication factor three.

How does delta-of-delta encoding get timestamps to about one bit?

With a fixed 10-second interval, the difference between consecutive timestamps is always 10 s, so the difference of differences is zero, encoded with a single 0 bit. Only jitter costs more bits. Values use XOR with the previous value: similar floats share most bits, so only the changed middle bits are stored.

What is a posting list?

For each label pair (service="api"), a sorted list of the ids of the series that have it. A query with several matchers intersects those lists, exactly as a search engine intersects the documents containing each word.

Why blocks of two hours?

Long enough to compress well and keep the number of objects manageable, short enough that ingester memory and WAL replay after a crash stay bounded. The compactor later merges them into larger blocks for efficient long-range reads.

How would you handle late or out-of-order samples?

Accept a small out-of-order window in the head (minutes), and reject or route older samples to a separate path. Most TSDBs assume nearly ordered data per series because that is what makes the compression work.

Cardinality

What is cardinality, and why is it the number one way to break a metrics system?

A series is a metric name plus a unique combination of label values. http_requests_total with labels for service (50), endpoint (100) and status (10) is up to 50,000 series. Add user_id with a million values and it becomes 50 billion. Every active series costs memory in the ingester and entries in the index, whether or not anyone queries it.

Cardinality explosions are usually accidental: a request id, a user id, an error message or a timestamp put into a label. They arrive in minutes with a deploy, and they hurt everyone who shares the ingesters.

  • Per-tenant and per-metric active-series limits at ingest, with clear errors, plus usage reporting and label guidance chosen
  • No limits; scale the cluster rejected
  • Drop or aggregate high-cardinality labels automatically situational: as an explicit, configured rule for known noisy labels

The answer: The gateway tracks active series per tenant and per metric (with an approximate counter such as HyperLogLog per window) and rejects new series beyond the limits while accepting samples for existing ones. Rejections return the metric and label responsible and raise an alert to the owning team. Query-side limits cap the series and samples one query can touch. Usage dashboards show each team its top metrics by series count, and guidance says what belongs in metrics (bounded labels) versus logs and traces (ids). High-cardinality analysis per request is the job of tracing or event analytics, not metrics.

A team really needs per-customer metrics for 10,000 customers. Allowed?

Ten thousand is bounded and probably fine with a raised limit and a cost conversation. A million customers would not be: use logs or an analytics store keyed by customer, and keep metrics at the level of plan or region.

How do you count active series per tenant cheaply?

Each ingester knows its own series; the gateway keeps an approximate distinct count of series hashes per tenant and metric over the last hour, which is accurate to a percent or two with a few KB of state. Exactness is not needed to enforce a limit.

How do you find which label caused an explosion?

Track series count per metric and per label name over time. A metric whose series count jumps after a deploy, with one label's distinct values growing fastest, is the culprit. Exposing this as a self-serve page saves the platform team from debugging every incident.

Should limits be soft or hard?

Both: a soft limit that warns the team early and a hard limit that rejects new series. The gap gives teams time to fix labels before anything is dropped.

A year of data at shrinking resolution

Someone opens a graph for the last 12 months. How do you avoid reading billions of samples?

A year of one series at 10-second resolution is 3.15 M samples; a graph is about 1,000 pixels wide. Reading every sample to draw a thousand points wastes three orders of magnitude of work, multiplied by every series in the query.

Averages alone are not enough to downsample: a 1-hour average hides a 30-second spike that the max would show, and averaging averages gives the wrong answer for counters and percentiles.

  • Rollups at 5 minutes and 1 hour storing min, max, sum and count; the query engine picks resolution from the range chosen
  • Keep raw data forever and compute on read rejected
  • Store only averages per hour rejected

The answer: Raw resolution is kept for 15 days, 5-minute rollups for 90 days, and 1-hour rollups for 13 months. Each rollup point stores min, max, sum and count, so averages, rates and peaks remain correct; latency distributions are stored as histogram buckets (counters), which can be summed across rollups and turned into percentiles at query time. The query engine chooses the coarsest resolution that still gives at least one point per pixel for the requested range and step, and splits long queries into day-aligned pieces that are cached.

Why can’t you average p99s across hosts or hours?

Percentiles do not compose: the average of two hosts’ p99 is not the p99 of their combined requests. Store the histogram buckets as counters, sum them across hosts and time, and compute the percentile at the end. That is why histograms, not pre-computed percentiles, are the right metric type.

A graph switches from raw to rollups and the line changes shape. Is that a bug?

It is the resolution changing, and the panel should show which one is in use. Using max for the line by default on long ranges keeps spikes visible, which is what people usually look for.

How do you roll up counters correctly?

Counters reset when a process restarts. Rollups store the increase within each interval (handling resets), not the raw counter value, so rates computed from rollups stay correct across restarts.

Can users still get raw data older than 15 days?

Not from the metrics system. If they need it, the samples can also be archived to cheap object storage in a columnar format for offline analysis. Keeping raw data hot for a year would multiply storage by about 30 for queries almost nobody runs.

Alerting that works when everything is broken

How do you make sure alerts fire, reach the right person, and do not bury them in noise?

Alerts are needed most during big incidents, exactly when the systems being watched, and possibly parts of the monitoring system, are failing. If alerting depends on the same storage and query path as dashboards, a monitoring outage silences every alert at once.

The opposite failure is noise: a database outage makes fifty services alert, on-call gets two hundred pages, and the one that names the cause is lost. Good alerting is as much about grouping and suppression as about evaluation.

  • Redundant rule evaluators on recent in-memory data, a deduplicating and grouping alert manager, routing by ownership labels, and an independent dead man’s switch chosen
  • Evaluate alerts from dashboards’ queries via the full query path rejected
  • Page on every threshold breach immediately rejected

The answer: Rules are sharded across evaluators, each shard run on two evaluators, querying only the last few minutes from ingester memory. Every alert has a "for" duration to ignore blips, and state per label set. Alerts go to a clustered alert manager that deduplicates the redundant copies, groups by service or cluster, applies silences and inhibitions (suppress symptom alerts while the cause is firing), and routes by owner labels to PagerDuty or Slack with retries. A heartbeat alert fires continuously into an external service that pages if it stops. Teams are encouraged to alert on user-facing symptoms (error rate, latency SLOs) rather than every cause.

An alert flaps between firing and resolved every few minutes. How do you calm it?

Add hysteresis: fire above 2 % for 5 minutes but resolve only below 1.5 % for 10 minutes, and keep a minimum notification interval. Often the real fix is a better rule, such as a burn-rate alert on an SLO, which is naturally smoother.

What is an SLO burn-rate alert?

Instead of alerting on a raw threshold, alert when the error budget is being consumed too fast: for example, at a rate that would exhaust a month’s budget in two days, measured over both a short and a long window. It pages for real user impact and ignores harmless blips.

Should the alert manager run in the same region as the systems it watches?

Run it clustered across zones, and have the dead man’s switch reach an external provider. For the most critical alerts, a second, minimal path (for example cloud-native alarms on a few key signals) gives independence from the whole stack.

Answering queries across shards and time

A query spans 100 ingesters and a year of blocks. How does it stay under a second?

Series are spread across ingesters by hash, so a query that aggregates over a label can touch every ingester. Time ranges older than a few hours live in many blocks in object storage, which is slow to read cold.

Most dashboard traffic is repetitive: the same panels, refreshed every 30 seconds, over ranges ending "now". Most of the work is identical to what was done 30 seconds ago.

  • Split by time into aligned pieces, cache each piece, push aggregation down to shards, read blocks through caching gateways chosen
  • Send every query to every shard and merge raw series centrally rejected

The answer: The query frontend aligns the time range to the step, splits it into day-sized pieces, and serves completed pieces from a results cache, computing only the most recent piece. Each piece is fanned out to the ingesters (for the last few hours) and to store gateways (for blocks), which use block indexes and caches to read only matching chunks. Where possible, aggregation is pushed down so shards return partial sums rather than raw series. Per-query limits on series, samples and time protect the fleet, and slow queries are logged with their cost so heavy dashboards can be fixed.

A sample arrives late for a time range already cached. Does the dashboard show stale data?

Only for pieces older than the out-of-order window. Do not cache the most recent piece (or cache it for seconds), and only cache pieces older than the late-arrival window, so late samples are reflected where they can still occur.

How do you protect the query fleet from one expensive dashboard?

Limits per query, per tenant concurrency limits, a queue per tenant in the frontend so one tenant cannot fill all workers, and cost reporting that names the panel. Recording rules can precompute expensive expressions that dashboards use constantly.

What are recording rules?

Queries that the system evaluates continuously and stores as new series, for example the per-service error rate. Dashboards and alerts then read the precomputed series instead of aggregating thousands of raw series on every refresh.

How do you keep query latency predictable for the last hour, the most common range?

Serve it from ingesters' memory, which is fast, and keep the head chunks small enough to scan quickly. Avoid caching the newest piece for long so dashboards stay live.

The theory behind it

  • Observability, operations and rollouts: Metrics, logs and traces; SLIs, SLOs and error budgets; alerting; safe deploys (canary, blue-green, feature flags); migrations without downtime; capacity planning; and how to answer "how would you know it is working".
  • Batch and stream processing: Stream versus batch, windowing and watermarks, late and out-of-order data, exactly-once counting, approximate algorithms, and the lambda/kappa argument in one paragraph.
  • Object storage and large files: Blob storage, pre-signed uploads, chunking and resumable transfer, multipart, erasure coding, storage tiers and lifecycle, plus how media pipelines are structured.

Related

Open in your browser to sign in

Google does not allow sign-in inside this app's built-in browser. Open this page in Safari and sign in there. The link opens this same page.

Tap the ⋯ or share button at the top or bottom of the screen, then Open in browser. Or copy the link and paste it into Safari.