Distributed Metrics Monitoring (Datadog / Prometheus scale)
Sharpened prompt. Design a metrics platform ingesting 100M data points/sec from 5M hosts and containers, storing 15 months, answering a 90-day dashboard query in under 2 seconds, and evaluating 500k alert rules every 30 seconds — while surviving the incident where one team ships a metric tagged with
request_id.
The last clause is the real problem. Metrics systems rarely die from volume; they die from cardinality, and the cardinality is generated by your own users, during an incident, when you need the system most.
1. Problem framing
Functional requirements
- Ingest counters, gauges, histograms, and distributions with arbitrary key-value tags.
- Query by metric name plus tag matchers, with aggregation over time and across series.
- Alerting: threshold, rate-of-change, and anomaly rules with per-rule evaluation windows.
- Retention tiers: raw for days, rolled up for months.
Non-functional requirements
| Property | Target | Consequence |
|---|---|---|
| Ingest | 100M points/sec, no drops | Agent-side aggregation is mandatory |
| Query | 90-day dashboard < 2s | Pre-computed rollups; raw data cannot be scanned |
| Freshness | Data queryable within 10s | Ingest path writes to a queryable memory buffer, not just to disk |
| Retention | 15 months | Tiered storage; raw data is deleted, not archived |
| Availability during incidents | Higher than the systems it monitors | Independent failure domain, separate cloud account/region |
Back-of-the-envelope
Points: 100M/sec × 86,400 = 8.6 trillion/day
Naive size: 16 bytes/point (8 ts + 8 value) = 138 TB/day -> impossible
Compressed: Gorilla gets ~1.37 bytes/point -> 11.8 TB/day, ~4.3 PB/15mo (rolled up: ~15%)
Series: 5M hosts × 200 metrics × 5 tag combos = 5 billion active series
Index: 5B series × 100 B of posting-list overhead = 500 GB of inverted index
Alerts: 500k rules / 30s = 16,600 evaluations/sec, each a mini-query
The compression number is the one to lead with: without Gorilla-class encoding this system does not exist, so it is not an optimisation, it is the foundation.
2. High-level architecture
Three things to call out on this diagram: the agent aggregates before the network (10× reduction), the ingester holds a queryable in-memory head block so recent data is fresh without waiting for a flush, and the alert evaluator is a client of the query layer rather than a parallel read path — one query engine, one set of semantics.
3. Component inventory
| Component | Concrete choice | Why this one |
|---|---|---|
| Agent | OpenTelemetry Collector / Vector | Pre-aggregates locally; a 10s flush turns 1000 requests into one histogram |
| Transport | Kafka partitioned by hash(series_id) | All points for a series land on one ingester, so compression state stays local |
| Ingester | Prometheus TSDB head / Cortex / Mimir | Append-only WAL plus a queryable in-memory block |
| Compression | Gorilla: delta-of-delta timestamps + XOR values | 12 bytes → ~1.37 bytes/point on real data |
| Index | Inverted index with Roaring bitmaps | Tag matching is set intersection; Roaring makes it fast and small |
| Block store | S3 with a local NVMe cache | Blocks are immutable, so object storage is ideal and 20× cheaper than EBS |
| Query engine | PromQL-compatible (Mimir / Thanos / VictoriaMetrics) | Standard semantics, and users already know it |
| Alerting | Sharded rule evaluator + Alertmanager for grouping | Dedup, grouping, and silencing are their own hard problem |
4. The toughest parts
4.1 Cardinality: the way this system actually dies
Why it's hard. A time series is uniquely identified by its name plus the full set of tag values. Cardinality is the product of tag value counts, not the sum. A developer adds user_id to a latency metric on a service with 10M users and creates 10M series from one metric line. Memory in the ingesters doubles, the index explodes, queries that used to touch 100 series now touch 10M, and the whole tenant's dashboards time out. This routinely happens during an incident, because that is when people add debugging tags.
http_requests{service, endpoint, method, status, host}
20 services × 50 endpoints × 5 methods × 8 statuses × 5000 hosts = 200M series
Add one tag: ... × request_id (unbounded) -> effectively infinite
Solution — three layers of defence, all necessary.
(a) Separate the index from the data. Do not store tags with every point. Assign each unique tag-set a series_id, keep tags → series_id in an inverted index, and store the samples as (series_id, timestamp, value). The index is where cardinality costs live, and it is now a bounded, separately scalable component.
Index (Roaring bitmaps, one posting list per tag pair):
service="checkout" -> {1, 5, 9, 12, ...}
status="500" -> {5, 12, 88, ...}
Query service=checkout AND status=500 -> bitmap AND -> {5, 12}
Roaring bitmaps matter here: intersecting two posting lists of millions of IDs is a handful of SIMD-friendly operations, and a dense list of 5M IDs compresses to a few hundred KB.
(b) Enforce a per-tenant cardinality budget at the gateway, and reject with a loud, specific error.
func (g *Gateway) Admit(tenant string, s Series) error {
// HyperLogLog per tenant: ~12 KB gives 2% error on billions of distinct series.
est := g.hll[tenant].Estimate()
if est > g.limits[tenant].MaxSeries {
// Reject the NEW series only; existing series keep flowing. Never drop
// a tenant's whole stream — that blinds them during an incident.
if !g.known[tenant].Contains(s.ID()) {
g.metrics.Inc("cardinality_rejected", tenant, s.Name)
return ErrCardinalityLimit{Metric: s.Name, Limit: g.limits[tenant].MaxSeries}
}
}
// Per-metric guard catches the single bad metric before it eats the budget.
if g.perMetric[tenant][s.Name].Estimate() > 100_000 {
return ErrMetricCardinality{Metric: s.Name, Tags: s.HighCardinalityTags()}
}
return nil
}
The subtlety worth stating: reject new series, keep existing ones. Dropping the whole tenant during their cardinality incident removes the observability they need to fix it.
(c) Give people somewhere else to put high-cardinality data. user_id and request_id are legitimate debugging needs — they just do not belong in a metric. Route them to traces (sampled, with exemplars linking metrics to traces) or to a wide-event store (ClickHouse). The senior framing: "Metrics are for aggregates with bounded dimensionality; anything unbounded belongs in traces or events. Cardinality limits aren't just protection, they're an interface contract."
4.2 Storing 8.6 trillion points/day: Gorilla compression
Why it's hard. At 16 bytes per point, 100M points/sec is 138 TB/day and 62 PB over 15 months. The storage bill alone kills the product. General-purpose compression (gzip, zstd) on raw arrays helps but is far too slow to run at ingest rate and does not exploit the structure.
Solution — exploit the two things that are always true of time series: timestamps are nearly regular, and consecutive values are nearly equal.
Timestamps: delta-of-delta. Samples arrive every 10s, so the delta is 10,000ms every time and the delta of deltas is 0. Encode 0 in a single bit.
t: 10:00:00 10:00:10 10:00:20 10:00:30
Δ: 10000 10000 10000
ΔΔ: 0 0 -> '0' bit each. 1 bit per timestamp.
Values: XOR with the previous value. A CPU gauge goes 0.71 → 0.72; the two doubles share nearly all their high bits, so the XOR is mostly zeros. Store the leading-zero count, the meaningful-bit count, and the meaningful bits.
def xor_encode(prev: float, cur: float, out: BitWriter):
x = struct.unpack('<Q', struct.pack('<d', prev))[0] ^ \
struct.unpack('<Q', struct.pack('<d', cur))[0]
if x == 0:
out.write_bit(0) # identical value: 1 bit total
return
out.write_bit(1)
lead, trail = clz(x), ctz(x)
if reuse_previous_window(lead, trail):
out.write_bit(0) # same window as last time: 2 bits + payload
out.write_bits(x >> trail, 64 - prev_lead - prev_trail)
else:
out.write_bit(1)
out.write_bits(lead, 5)
out.write_bits(64 - lead - trail, 6)
out.write_bits(x >> trail, 64 - lead - trail)
Facebook's Gorilla paper reports ~1.37 bytes per point on production data — a 12× reduction, turning 138 TB/day into ~12 TB/day. That single change is the difference between a viable product and a bankrupt one.
Two operational consequences to mention: chunks must be immutable and closed (typically 2 hours or 120 samples) because the encoding is a stream, and the encoding is not seekable — you decompress a chunk from its start, which is why chunks stay small.
For counters, delta encoding plus the fact that they only increase compresses even further; for histograms, store the bucket boundaries once in the index and only the counts per chunk.
4.3 Ingesting 100M points/sec without dropping any
Why it's hard. 100M points/sec of individual HTTP requests is impossible — the syscall and TLS overhead alone would need tens of thousands of cores. Worse, ingest load spikes exactly during incidents, when every service starts emitting error metrics and retry counters, and dropping data then is the worst possible time.
Solution — aggregate early, batch hard, and shard by series.
- Agent-side pre-aggregation. The agent maintains local counters and histograms and flushes every 10 seconds. A service handling 10,000 requests/sec emits one histogram per 10s, not 100,000 points. This is a 10–1000× reduction and it happens before any network hop.
- Sketches, not raw values, for percentiles. Never ship raw latencies to compute P99 centrally. Use DDSketch (relative-error guarantee, mergeable) or t-digest. A DDSketch with 1% relative error is ~2 KB and merges across hosts, so a fleet-wide P99 is the merge of 5,000 sketches rather than a sort of 50M values. Note that averaging per-host P99s is mathematically wrong — mergeable sketches are the only correct way to get a fleet percentile, and saying so is a strong signal.
- Partition Kafka by
hash(series_id)so every point for a series reaches the same ingester, keeping the compression stream and the head chunk local. Partitioning by host would scatter one series across ingesters and break both. - Backpressure with priority. When ingesters saturate, shed low-value data first: debug-level metrics before SLO metrics. Publish the shed rate as a metric itself so tenants can see what they lost.
# OTel Collector: the 10s aggregation window is the highest-leverage config here.
processors:
batch: {timeout: 10s, send_batch_size: 8192}
cumulativetodelta: {}
filter/drop_debug:
metrics: {exclude: {match_type: regexp, metric_names: ["debug\\..*"]}}
exporters:
kafka: {brokers: [...], encoding: otlp_proto, producer: {compression: snappy}}
4.4 A 90-day query in under 2 seconds
Why it's hard. 90 days of 10-second data is 777,600 points per series. A dashboard panel over 5,000 hosts is 3.9 billion points to read, decompress, and aggregate — for a chart that is 1,200 pixels wide and can therefore display at most 1,200 values. You are doing a million times more work than the output justifies.
Solution — pre-computed rollups plus query planning that picks the right resolution.
Retention and resolution tiers:
raw 10s -> 7 days (incident debugging)
1m -> 30 days (weekly analysis)
1h -> 13 months (capacity planning)
1d -> 15 months (year-over-year)
Each rollup stores count, sum, min, max plus a mergeable sketch for percentiles — never a pre-computed P99, which cannot be re-aggregated across a different grouping.
// Query planner: choose the coarsest resolution that still fills the pixels.
func pickResolution(from, to time.Time, pixels int) time.Duration {
span := to.Sub(from)
want := span / time.Duration(pixels) // seconds per pixel
for _, r := range []time.Duration{10*time.Second, time.Minute, time.Hour, 24*time.Hour} {
if r >= want { return r }
}
return 24 * time.Hour
}
// 90d over 1200px -> 6480s/px -> the 1h tier. 3.9B points becomes 10.8M. 360× less work.
Then layer the standard query-frontend tricks: split a long query into per-day sub-queries executed in parallel, cache completed days (immutable, so cacheable forever) and re-run only the partial current day, and enforce limits — max series touched, max samples scanned, max duration — so one bad query cannot starve the cluster. Return the limit error with the offending matcher named; users learn quickly when the error is specific.
4.5 Evaluating 500k alert rules every 30 seconds
Why it's hard. Each rule is a query. 500k queries per 30 seconds is 16,600 QPS of analytical queries, all firing at the same instant if you schedule them naively, and every one competing with human dashboard traffic. Meanwhile a single flapping rule can generate thousands of pages.
Solution.
(a) Shard rules by hash and stagger within the interval. Rule r evaluates at offset hash(r) % 30s, spreading 500k evaluations evenly instead of stacking them on the 30-second boundary. Same deterministic-jitter trick as the job scheduler.
(b) Deduplicate the underlying queries. Hundreds of rules watch rate(http_errors[5m]) with different thresholds. Evaluate the expression once per group and apply many predicates to the result.
(c) Alert on symptoms, and require for durations. A rule that fires the instant a threshold is crossed will flap. for: 5m requires the condition to hold continuously, which eliminates the large majority of noise.
groups:
- name: slo
interval: 30s
rules:
- alert: CheckoutErrorBudgetBurn
# Multi-window burn rate: fast burn (1h) confirmed by a longer window (5m)
# to avoid paging on a 30-second blip.
expr: |
(slo:error_ratio:rate1h{svc="checkout"} > 14.4 * 0.001)
and
(slo:error_ratio:rate5m{svc="checkout"} > 14.4 * 0.001)
for: 2m
labels: {severity: page}
annotations: {summary: "Burning 30-day budget in ~2 days"}
Multi-window multi-burn-rate alerting is worth knowing by name: it pages fast on real outages and stays quiet on blips, and it is the current best practice for SLO alerting.
(d) Group and inhibit downstream. A rack losing power fires 200 host-down alerts. Alertmanager groups by cluster, alertname into one notification and inhibits dependent alerts when a parent alert is firing. Without this the alerting system becomes noise and people mute it, which is worse than having no alerts.
4.6 Late, out-of-order, and skewed samples
Why it's hard. Prometheus's TSDB traditionally requires strictly increasing timestamps per series and rejects out-of-order samples outright. But an agent buffers during a network partition and delivers 10 minutes late; a container's clock is 3 seconds off; a Kafka partition lags and delivers behind a fresher one. Rejecting all of it means silent gaps exactly around the interesting moments.
Solution.
- Timestamp at the source, in the agent, not at the ingester. Ingest-time stamping makes a delayed batch look like a spike of simultaneous samples — a lie that shows up in every chart.
- Accept a bounded out-of-order window (Prometheus's
out_of_order_time_window, Mimir's OOO head block): keep a separate head block for late samples and merge at query time. One hour of allowed lateness covers most agent buffering. - Discipline clocks with chrony and monitor skew as a metric. A host more than 5 seconds off gets flagged; its data is still accepted but marked suspect.
- Beyond the window, route to a reconciliation path that rewrites the rollups for the affected period, and mark the query result as "revised" so a dashboard viewed twice showing different numbers has an explanation.
4.7 Monitoring the monitoring system
Why it's hard. If the metrics platform runs in the same region, same Kubernetes cluster, and same cloud account as the systems it monitors, then a regional outage takes out both — and you are blind precisely during the largest incident you will ever have. This is a genuine, repeated, industry-wide failure.
Solution — an independent failure domain plus a minimal external watchdog.
- Run the monitoring stack in a separate account and region from the primary workload, with no shared control plane.
- Run a tiny secondary monitor (a "meta-monitor") on entirely different infrastructure — a different cloud, or a hosted service — whose only job is to answer "is the main monitoring system ingesting and alerting?" It watches a handful of heartbeat metrics, not everything.
- Dead man's switch: a rule that fires continuously while healthy, wired to an external service that pages when the stream stops. Absence of signal is the failure you cannot otherwise detect.
- Make the alerting path degrade independently: if the query layer is down, alert evaluation should fall back to a simplified path over the ingesters directly rather than going silent.
5. What breaks first
| Event | First failure | Mitigation |
|---|---|---|
A team adds request_id as a tag | Ingester memory, then index size | Per-tenant and per-metric cardinality limits at the gateway; reject new series only |
| Region-wide incident | Ingest doubles from error metrics | Agent aggregation absorbs most of it; priority-based shedding for the rest |
| Someone queries 15 months at 10s resolution | Querier OOM, cluster-wide impact | Query limits on series/samples/duration; automatic resolution selection |
| Ingester crash | Last 2h of head block at risk | WAL replay on restart; replication factor 3 with quorum writes |
| S3 throttling on block upload | Flush backpressure into ingester memory | Exponential backoff, local NVMe spill, prefix-sharded object keys |
| 500k rules all evaluate at :00 | Query layer saturates every 30s | Deterministic per-rule offset within the interval |
| Monitoring shares fate with production | Blind during the biggest incident | Separate account/region, external dead man's switch |
6. Cheat sheet
- Cardinality is the enemy. Series = product of tag values. Index separate from data, per-tenant HLL budgets, and a documented "put it in traces instead" escape hatch.
- Gorilla compression: delta-of-delta timestamps + XOR values ≈ 1.37 bytes/point, 12× smaller. Non-negotiable at this scale.
- Ingest: agent aggregates over 10s; DDSketch/t-digest for percentiles because averaging P99s is wrong; Kafka partitioned by series hash.
- Query: rollup tiers (10s/1m/1h/1d), planner picks resolution from pixel count, split + cache immutable days, hard limits.
- Alerting: shard and stagger rules, dedupe shared expressions, multi-window burn rates, group and inhibit.
- Late data: stamp at the source, bounded out-of-order head block, monitor clock skew as a first-class metric.
- The one-liner: "This is two databases pretending to be one — an inverted index over tag sets, and a compressed append-only column store of floats. Every hard problem lands on one side or the other, and cardinality is what kills the index."