On this page
System Design — Metrics Monitoring & Alerting
Last reviewed 11 Sept 2026
Part of the system design series. See the framework and building blocks first if you haven’t.
1. Requirements
Functional
- Services emit numeric metrics (counters, gauges, histograms) tagged with labels (
service,region,status_code). - Query metrics over time ranges with aggregation (
avg,p99,rate()) via a query language. - Define alert rules over queries; notify on-call when a threshold is breached for a sustained duration.
- Dashboards for humans, ad hoc queries for debugging an incident.
Non-functional
- Ingest at very high cardinality and volume — potentially millions of unique time series across a large fleet.
- Write-heavy, append-only, time-ordered — a fundamentally different access pattern from an OLTP database.
- Query latency low enough for a human staring at a dashboard during an incident, and for an alert-evaluation loop running every 10-30 seconds.
- The monitoring system must stay up especially when the rest of the infrastructure is on fire — it cannot depend on the very systems it’s watching.
2. Where it sits / high-level architecture
flowchart TD S1[Service instance] -->|expose /metrics| SC[Scraper - pull model] S2[Service instance] -->|expose /metrics| SC SD[Service discovery] -.tells scraper what to scrape.-> SC SC --> TSDB[(Local time-series DB<br/>per region/cluster)] TSDB -->|remote_write| LTS[(Long-term store<br/>object storage + compactor)] TSDB --> RULE[Alert rule evaluator] RULE -->|firing| AM[Alertmanager<br/>dedupe/group/route] AM --> PD[PagerDuty / Slack / on-call] Q[Dashboard / ad hoc query] -->|PromQL| GLOBAL[Global query layer] GLOBAL --> TSDB GLOBAL --> LTS
- Pull (scrape) model: a central scraper periodically hits each instrumented service’s
/metricsendpoint, discovering targets dynamically via service discovery (Kubernetes API, Consul). This is Prometheus’s defining choice — it inverts the usual “services push to a collector” pattern. - Local TSDB per region/cluster holds recent, high-resolution data close to where it’s generated, keeping the hot write path local and cheap.
- Remote-write to a long-term, globally queryable store (object storage-backed) handles retention beyond what local disk can hold and gives a single query surface across regions.
- Alert rule evaluator runs continuously against the local TSDB (not the slower long-term store) so alerting latency stays low, and hands firing alerts to a dedicated Alertmanager-style component for dedup/grouping/routing — deliberately a separate concern from evaluation.
3. Ingestion model and storage trade-offs
| Choice | What you get | What you give up |
|---|---|---|
| Pull (Prometheus-style scrape) | Central control over scrape frequency/targets; trivially know if a target is down (scrape fails); no risk of a buggy service DDoSing the collector | Scraper needs network reachability to every target; awkward for short-lived batch jobs (use a pushgateway) |
| Push (StatsD/Datadog agent-style) | Works naturally for ephemeral/batch jobs and anything the collector can’t reach | Collector must handle arbitrary incoming load and is a more attractive target for a misbehaving client to overwhelm |
| Row-oriented storage | Simple to implement | Terrible compression and scan performance for “give me one metric’s values over 30 days” queries |
| Columnar / time-series-native storage with delta-of-delta + XOR compression (Gorilla-style, from Facebook’s 2015 paper) | ~10x+ compression on typical metric data — timestamps compress via delta-of-delta encoding, float values via XOR-based bit-packing since consecutive samples are usually close in value | More specialized storage engine, not a general-purpose DB |
| Downsampling old data (5-min, then 1-hour resolution as data ages) | Bounded long-term storage cost despite unbounded retention | Old incident investigations lose fine-grained resolution — acceptable trade, nobody needs second-level detail from 6 months ago |
High cardinality is the recurring failure mode to name: a label like user_id or a raw request path with embedded IDs turns one logical metric into millions of unique time series, blowing up memory and index size. The fix is architectural, not just operational — reject or aggregate away high-cardinality labels before they hit the TSDB, don’t try to tune your way out after the fact.
4. Deep dive — the two hardest parts
a) Making “local, fast alerting” and “global, durable, long-term storage” coexist without either slowing the other down. The naive design — one big central time-series database everything writes to and every alert rule queries — creates a single point of failure and a latency/scale bottleneck exactly where you can least afford one (alert evaluation during an incident). The production pattern (Prometheus + Thanos/Cortex/Mimir/M3) splits this cleanly: each region/cluster runs its own local Prometheus doing its own scraping and its own alert-rule evaluation against local data only — so alerting never depends on a remote system. A sidecar or remote-write path asynchronously ships data to a separate long-term/global-query layer (object storage plus a compactor/querier), which is what dashboards and cross-region ad hoc queries hit, decoupled from the alerting hot path entirely.
b) Alert rule evaluation without flapping. A naive “value crossed threshold, fire immediately” rule flaps constantly on noisy metrics (a latency spike that resolves in one scrape interval shouldn’t page anyone). The real answer has two parts: (1) require the condition to hold for a sustained duration (e.g. “p99 latency > 500ms for 3 consecutive scrapes/minutes,” not one sample), and (2) separate evaluation from notification — Alertmanager-style components take a stream of firing/resolved alert states and apply deduplication (don’t send the same alert twice), grouping (bundle related alerts from one root cause into one notification), inhibition (suppress downstream alerts when a known upstream cause is already firing — e.g. don’t page for every service’s error rate spike if the shared database is already flagged down), and routing (which team, which channel, escalation policy) before anything reaches a human.
5. What real systems do today
- Prometheus is the dominant open-source model: pull-based scraping, its own on-disk time-series storage, PromQL for both dashboards and alert rules, and dynamic service discovery (notably Kubernetes-native) rather than static target lists.
- Long-term storage extensions — Thanos, Cortex/Mimir, and M3DB — are the standard answer to Prometheus’s single-node local-disk retention limits. Thanos keeps Prometheus itself as the ingestion/alerting layer and bolts on a sidecar that uploads TSDB blocks to object storage (S3/GCS) plus a separate Store Gateway and Compactor for historical queries and downsampling (commonly to 5-minute and 1-hour resolutions). Cortex/Mimir instead centralize around a remote-write ingestion path into their own storage backend. Both approaches explicitly decouple “fast local alerting” from “cheap, durable, long-range storage.”
- Datadog and similar commercial platforms lean toward a push/agent model rather than pull, trading the pull model’s “is the target even reachable” simplicity for easier coverage of ephemeral and non-HTTP-exposed workloads.
- Compression technique lineage traces to Facebook’s Gorilla paper (delta-of-delta timestamps, XOR’d floats), which is the widely cited justification for why purpose-built time-series storage engines outperform general-purpose databases by a large factor for this workload.
- Current (2025-2026) engineering guidance on choosing between Thanos/Cortex/Mimir/M3 for long-term storage converges on: Thanos for teams that want to keep Prometheus as-is and bolt on global/long-term query capability with minimal architecture change; Cortex/Mimir for teams wanting a more centralized, horizontally-scaled ingestion path from the start.
6. Scaling & failure
| Bottleneck | Fix | New cost |
|---|---|---|
| Single Prometheus instance can’t hold the metric volume of a large fleet | Shard scrape targets across multiple Prometheus instances (functional sharding, e.g. by team/service) | Cross-shard queries (a dashboard spanning multiple shards) need the global query layer, not a single instance |
| Local disk retention limits how far back you can query | Remote-write / sidecar upload to object storage with a compactor doing downsampling | Slightly higher latency and cost for old-data queries; acceptable since incident investigation rarely needs full resolution from months back |
| High-cardinality labels blow up memory/index on a single node | Reject/aggregate high-cardinality labels at the instrumentation or scrape-relabel layer before ingestion | Loses per-value granularity for that label — must be a deliberate trade agreed with the team emitting the metric, not a silent drop |
| Alert notification storm during a large correlated outage | Alertmanager grouping + inhibition rules (suppress downstream alerts when a known root-cause alert is already firing) | Requires maintaining an explicit dependency map between alerts, which needs upkeep as the system evolves |
What happens when the long-term/global storage layer dies
- Because alert evaluation runs against local per-region storage, not the long-term store, active alerting keeps working through a long-term-storage outage — this is the entire point of the split architecture, and it’s the answer to give first.
- What breaks: cross-region dashboards, historical queries beyond local retention, and any ad hoc investigation reaching further back than the local TSDB’s window. During the outage, local Prometheus instances keep scraping and buffering normally (subject to local disk limits) and typically backfill the remote store once it recovers.
- What happens when local scraping/storage itself dies (the more serious case): that region loses both alerting coverage and recent data for the outage window — this is the actual single point of failure to call out, and the standard mitigation is running local Prometheus/TSDB instances redundantly (e.g. two independently scraping the same targets) so losing one doesn’t blind that region, plus ensuring the monitoring stack’s own infrastructure (its compute, its own alerting-about-itself) is provisioned separately from the services it watches, so a shared-infra outage doesn’t take out your ability to see the outage.
Interview follow-ups
- “Why pull instead of push?” — Central control over scrape targets/frequency, trivial dead-target detection (scrape failure itself is a signal), no risk of a buggy service flooding the collector; trade-off is needing network reachability to every target and an awkward story for short-lived batch jobs (pushgateway pattern).
- “How do you keep alerting fast when your metrics volume is huge?” — Alert rules evaluate against local, regional storage, not the global long-term store — decouple the hot alerting path from the slower durable/global query path entirely.
- “What’s high cardinality and why is it dangerous?” — A label with unbounded unique values (user_id, raw path with embedded IDs) multiplies one logical metric into millions of time series, blowing up memory and index size; fix it before ingestion, not after.
- “How do you stop an alert from flapping on a one-sample noise spike?” — Require the condition to hold for a sustained duration across multiple evaluation cycles, not a single breach.
- “A big outage causes 200 services to all fire error-rate alerts at once. What happens to your on-call?” — Alertmanager-style grouping and inhibition rules collapse correlated alerts into one notification tied to the likely root cause, instead of paging 200 times.
- “What compression technique makes time-series storage so much cheaper than a generic database for this data?” — Gorilla-style delta-of-delta timestamp encoding plus XOR-based float compression, exploiting that consecutive samples are usually close in time and value.
- “Your long-term storage (Thanos/Cortex) goes down. What’s the actual user impact?” — Real-time alerting is unaffected (it’s local); historical/cross-region dashboards and old-data queries degrade until it recovers.
Sources: Designing a Metrics and Monitoring System: Prometheus at Scale — DEV Community · Scaling Prometheus with Thanos: Long-Term Storage, HA, and Global Queries — Exoscale · What end-users want out of Prometheus remote storage: M3 vs Thanos — Chronosphere · Prometheus Storage Scaling 2026: Thanos vs Cortex vs Mimir · Prometheus Monitoring: How It Works, Key Features & Implementation Guide — CubeAPM · Design a Metrics & Monitoring System (like Datadog/Prometheus) — DesignGurus