Design a Metrics and Monitoring System
Run this for someone else. You hold the answers; they do not. Read the prompt, keep the clock, and use the probes below when an answer is thin. Do not show them this page.
Open with this
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. Take a couple of minutes on requirements, then we will do some numbers, then the design. I will interrupt to keep us moving.
The clock
- 4 min: functional requirements and scope
- 4 min: non-functional requirements, with numbers
- 5 min: back-of-envelope estimates
- 16 min: high-level design and one or two flows
- 16 min: deep dives and the close
Move them on out loud when a section overruns. The commonest failure is spending twenty minutes on requirements and never reaching a deep dive, and preventing that is your job as much as theirs.
Requirements · 8 min
Listen for: a scoped set of capabilities, an explicit out-of-scope list, and numeric targets rather than adjectives. Prompt with “what are you not building?” if they never scope, and “what number would make that requirement real?” if they say “fast” or “highly available”.
Functional (6)
- 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 (6)
- 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.
Estimates · 5 min
Ask for two or three numbers, not all of them. What matters is whether they state assumptions, round sensibly, and say what the number implies. Push once with “where did that come from?”
The numbers (6)
- 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.
High-level design · 16 min
Let them draw. Interrupt only to ask what backs a component or what a box actually does. Then pick one flow below and ask them to walk it end to end.
Components (14)
- 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.
Flows to ask them to walk (5)
- 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.
- The agent sends a compressed batch of samples every 10 seconds.
- The gateway validates labels and enforces the tenant’s rate and series limits.
- Accepted samples are written to the ingest log, partitioned by a hash of the series’ labels.
- An ingester appends each sample to its series’ in-memory chunk and to its write-ahead log.
- Every two hours the ingester cuts an immutable block and uploads it to object storage.
- The compactor merges blocks, drops replica duplicates and writes rollups.
- 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.
- A panel asks for p99 latency by region over the last 6 hours.
- The cache returns everything except the newest step; only that piece is computed.
- The query engine uses the inverted index to find the matching series.
- Recent hours are read from ingesters; older data from blocks, at a resolution chosen from the range.
- Partial results are merged and aggregated, and replica duplicates removed.
- 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.
- A deploy adds a user_id label to a request counter.
- The gateway sees the tenant’s active series approaching its limit and rejects new series for that metric.
- Without limits, ingesters would run out of memory and the index would balloon.
- Queries touching the metric would scan millions of series.
- 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.
- One ingester dies with two hours of unflushed data in memory.
- A replacement starts, replays its write-ahead log, then resumes consuming from its last committed offset.
- Meta-monitoring notices the ingester was down and records the event.
- The compactor later removes the duplicate copies in uploaded blocks.
- 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.
- Every 30 seconds the rule evaluator runs "5xx rate above 2 % for 5 minutes" for each service.
- The condition holds; the alert goes pending, then firing after five minutes.
- Firing alerts are sent to the alert manager, which deduplicates, groups and routes them.
- The owning team is paged.
- A heartbeat alert that always fires proves the path works; its absence pages someone.
Deep dives · 16 min
Pick two. Ask the headline question, let them answer, then use the follow-ups. The follow-ups are where the level gets decided, so leave time for at least three of them.
Push or pull
Ask: Should the monitoring system scrape targets, or should agents push to it?
Good answers name: Agents on each host scrape local targets (pull) and push batches to a central ingest (push), Central pull (Prometheus servers scraping everything), Applications push directly to the backend.
Our pick: 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
Ask: Ten million samples a second. Why not just insert rows into a database?
Good answers name: Purpose-built TSDB: per-series compressed chunks in memory, WAL, immutable blocks with an inverted label index, in object storage, Wide-column store (one row per series per time bucket), Relational database with a timestamp index.
Our pick: 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
Ask: What is cardinality, and why is it the number one way to break a metrics system?
Good answers name: Per-tenant and per-metric active-series limits at ingest, with clear errors, plus usage reporting and label guidance, No limits; scale the cluster, Drop or aggregate high-cardinality labels automatically.
Our pick: 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
Ask: Someone opens a graph for the last 12 months. How do you avoid reading billions of samples?
Good answers name: Rollups at 5 minutes and 1 hour storing min, max, sum and count; the query engine picks resolution from the range, Keep raw data forever and compute on read, Store only averages per hour.
Our pick: 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
Ask: How do you make sure alerts fire, reach the right person, and do not bury them in noise?
Good answers name: 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, Evaluate alerts from dashboards’ queries via the full query path, Page on every threshold breach immediately.
Our pick: 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
Ask: A query spans 100 ingesters and a year of blocks. How does it stay under a second?
Good answers name: Split by time into aligned pieces, cache each piece, push aggregation down to shards, read blocks through caching gateways, Send every query to every shard and merge raw series centrally.
Our pick: 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.
Close · 5 min
Ask what breaks first at ten times the load, and what they would build next. Then give them your read: one thing that was strong, one thing that was missing, one thing to practise. Be specific; “good job” helps nobody.