Why This Matters#
A time series database is not a database choice. It is a cardinality budget with a storage engine attached. Every TSDB — Prometheus, InfluxDB, TimescaleDB, M3, VictoriaMetrics — is fast at the same thing: appending timestamped numbers to a known series and reading a contiguous time range back. Every one of them dies the same way: somebody adds a label with a million distinct values, the in-memory index explodes, and the monitoring system goes dark exactly when the incident it was supposed to catch begins.
That is why this technology keeps showing up in Staff loops. "Design a metrics platform", "Design ad click aggregation", "Design an IoT telemetry pipeline", and "How do you monitor this system?" are all TSDB questions in disguise. The L5 candidate names Prometheus and draws a Grafana box. The L6 candidate says "active series is my capacity unit, not bytes; I'll cap cardinality per team at 500K series and reject user_id as a label." The L7 candidate asks who pays for the series — and whether observability should be a chargeback platform or a free-for-all that bankrupts itself in year two.
The gap between levels is not knowing what delta-of-delta encoding is. It is knowing that the write path is almost never the bottleneck — the index is, that percentiles cannot be averaged, that retention is a product decision with a dollar figure, and that the monitoring system must fail independently of the thing it monitors.
The L5 → L6 → L7 Contrast#
| Behavior | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| First move | "Prometheus + Grafana" | "What is the active-series count and churn rate? That sizes everything." | "Is this a platform the org runs, or does each team own its own stack? That decides the next 3 years of cost." |
| Capacity unit | GB of disk | Active series × scrape interval → samples/sec; memory per series (~3–8KB) | $ per million active series per month, chargeable to the emitting team |
| Cardinality | "Add labels for filtering" | "Labels are bounded enums. IDs go to logs or traces, never labels." | Writes the org-wide label policy with enforcement at ingest and a quota exception process |
| Retention | "Keep 30 days" | Tiered: raw 15s × 15d, 5m rollups × 90d, 1h rollups × 2y | Retention is priced per tier and signed off by product/finance, not defaulted |
| Failure | "Run two Prometheus replicas" | Meta-monitoring: the monitoring system is watched by something outside its blast radius | Observability is a separate failure domain with its own error budget; correlated-failure review across regions |
| Build vs buy | "Use Datadog" or "self-host" | Compares ingest cost, query needs, and ops load | Models the crossover point (~1–5M series) where vendor per-series pricing exceeds a 3–4 engineer platform team |
Why "Capacity unit" separates levels
A Senior engineer sizes a TSDB like a database: bytes on disk. Disk is the cheapest part. Gorilla-style compression gets a sample to ~1.3–2 bytes, so 1M series at a 15s scrape is ~67K samples/sec, ~5.8B samples/day, ~10GB/day. Nobody gets paged for 10GB/day. What pages people is the head index: every active series costs memory for its label set, postings entries, and the open compression chunk. Prometheus typically burns 3–8KB of RAM per active series, so 10M series is 30–80GB of RAM on one box before a single query runs. The Staff candidate sizes by active series and churn, because that is the resource that runs out first.
Why "Cardinality" separates levels
Adding user_id to a request-latency metric feels like good observability. It turns 20 series (endpoint × status) into 20 × 50M. The L5 answer is not wrong about the need — per-user debugging is real — it is wrong about the tool. The Staff answer routes high-cardinality dimensions to logs, traces, or an OLAP store (ClickHouse, Druid) and keeps the TSDB for aggregate health. The Principal answer makes that routing a policy with an enforcement point, because one team's label mistake takes down alerting for all 200 teams sharing the cluster.
The 60-Second Pitch#
"For operational metrics I'd use a Prometheus-compatible TSDB — Prometheus at the edge for scraping and short-term alerting, remote-writing into a horizontally scalable long-term store like Mimir, Thanos, or VictoriaMetrics. The data is append-only, time-ordered, and numeric, so Gorilla-style compression gets us to under 2 bytes per sample — about 8–12× smaller than raw. The thing I'll actually design around is cardinality: I'm budgeting active series, not disk, and I'll enforce label limits at ingest. Retention is tiered — 15 days raw, 90 days at 5-minute rollups, 2 years hourly. If the question is business analytics with joins and ad-hoc slicing by customer, I'd switch to TimescaleDB or ClickHouse instead — a TSDB is the wrong tool for high-cardinality questions."
That pitch hits the four things interviewers listen for: a named engine with a reason, a compression number, the cardinality constraint, and the explicit boundary of when not to use it.
The Three Intents#
TSDB questions hide three different workloads. They lead to incompatible designs.
| Intent | Constraint | Strategy | Failure Mode | Correctness Bar |
|---|---|---|---|---|
| Operational monitoring & alerting | Freshness < 30s, must survive the incident it reports on | Pull-based Prometheus per cluster, local alerting, remote-write to long-term store | Monitoring dies with the cluster; alert storms | Approximate is fine; missing data is not |
| Business / product analytics over time | High cardinality (per customer, per SKU), ad-hoc SQL, joins | TimescaleDB, ClickHouse, Druid; columnar with rollups | Treating Prometheus as an analytics DB → OOM | Exact counts, auditable, reproducible |
| IoT / sensor telemetry | Millions of devices, out-of-order and late data, long retention | Push ingest via Kafka → InfluxDB 3 / TimescaleDB / VictoriaMetrics, device-bucketed | Per-device series churn; late data rejected | Every reading stored; gaps are visible |
🎯 Staff Move: "I'll assume operational monitoring first — that's where freshness, cardinality, and 'the monitor must outlive the outage' collide. If this is really billing-grade usage metering, I'd move it off the TSDB entirely onto an exactly-once pipeline, because a TSDB is allowed to drop samples and a bill is not."
The Staff Positions#
| Position | Rationale |
|---|---|
| Active series is the capacity unit | Memory for the head index runs out long before disk; size by series and churn, then derive samples/sec. |
| Labels are bounded enums; IDs never become labels | Cardinality is multiplicative. One unbounded label turns 1K series into 1B. |
| Alert at the edge, store in the core | Alerting runs on the Prometheus closest to the target so a WAN or central-store outage cannot silence pages. |
| Tier retention by resolution | Raw data older than ~2 weeks is almost never queried at raw resolution; rollups cut storage 20–240×. |
| Histograms, not averaged percentiles | p99 of p99s is meaningless. Store bucket counts and aggregate buckets. |
| The monitoring system is its own failure domain | Meta-monitoring from outside the blast radius, e.g. a dead-man's-switch alert that fires when the heartbeat stops. |
| Metering and billing do not live in a TSDB | TSDBs drop, deduplicate, and downsample by design; money needs exactly-once and audit. |
Architecture & Internals#
Only four internals change design decisions: the compression scheme, the head block + WAL, the inverted index, and block compaction. Everything else is implementation detail.
The Data Model — Series, Labels, Samples#
A series is identified by a metric name plus a sorted set of label key-value pairs: http_requests_total{service="checkout", method="POST", status="500"}. A sample is a (timestamp_ms int64, value float64) pair. Raw, that is 16 bytes. The storage engine's whole job is to make it ~1.5.
| Engine | Series identity | Schema | Query language |
|---|---|---|---|
| Prometheus / Mimir / Thanos / VictoriaMetrics | metric name + labels | Schemaless, labels only | PromQL (MetricsQL for VM) |
| InfluxDB 1.x/2.x | measurement + tag set; fields hold values | Tags indexed, fields not | InfluxQL / Flux |
| InfluxDB 3 | table + tag columns | Columnar (Arrow + Parquet) | SQL / InfluxQL |
| TimescaleDB | any PostgreSQL row with a time column | Full relational schema | SQL with time_bucket() |
| M3DB | ID + tags | Schemaless | PromQL / M3 query |
Gorilla Compression — Why Samples Cost ~1.4 Bytes#
Facebook's 2015 Gorilla paper introduced the encoding that nearly every modern TSDB copies. Two observations drive it:
- Timestamps are regular. Scrapes arrive every 15s, so the delta between timestamps is ~15000ms, and the delta of the delta is almost always 0. Store delta-of-delta with a variable-length prefix: a
0bit when it is zero. That is 1 bit per timestamp in the common case. - Consecutive values are similar. XOR a float with the previous float; most leading and trailing bits are zero. Store only the meaningful middle bits. An unchanged value costs 1 bit.
timestamps: 1000, 1015, 1030, 1045, 1061
deltas: 15, 15, 15, 16
delta-of-delta: 0, 0, 1 -> '0' '0' '10'+bits
values: 42.0, 42.0, 42.5, 42.5
xor w/ prev: 0, X, 0 -> '0' 'meaningful bits' '0'
Gorilla paper result: ~1.37 bytes/sample average vs 16 raw (~12x)
Prometheus in practice: ~1.3-2 bytes/sample on disk
Design consequences of this encoding:
- Counters compress better than gauges. A monotonically increasing counter with a stable rate has near-constant deltas.
- Jittery scrape intervals cost bytes. Irregular timestamps break the 1-bit path.
- High-precision random floats compress badly (~8+ bytes). Don't push sensor noise at 15 decimal places.
- Chunks are append-only, typically ~120 samples each. You cannot cheaply update a sample in the middle — updates and deletes are expensive in every TSDB.
The Write Path — Head Block, WAL, Blocks#
- Append to WAL for crash recovery. Replaying a large WAL is why a Prometheus with 10M series can take 5–20 minutes to restart — an operational fact that belongs in your failure-mode discussion.
- Append to the open chunk of the series in memory. New series also update the in-memory inverted index (label value → sorted list of series IDs, called postings).
- Every 2 hours the head is cut into an immutable on-disk block: chunks + index + tombstones.
- Compaction merges 2h blocks into larger blocks (in Prometheus, up to 31 days or 10% of the retention window), which shrinks the index and speeds long-range queries.
- Retention deletes whole blocks. There is no row-level delete in the hot path — deletes are tombstones applied at compaction.
🎯 Staff Insight: "The head block is the entire operational story. Memory scales with active series, restart time scales with WAL size, and series churn — pods coming and going with new label values — inflates the index even if the steady-state series count looks fine. I'd alert on
prometheus_tsdb_head_seriesand on churn rate, not on disk."
The Read Path — Postings Intersection#
A PromQL selector like http_requests_total{service="checkout", status=~"5.."} resolves by:
- Look up the postings list for
__name__="http_requests_total", forservice="checkout", and the union of postings for everystatusvalue matching5... - Intersect the sorted lists → matching series IDs.
- For each series, decode only the chunks overlapping the query time range.
- Apply functions (
rate,sum by,histogram_quantile) in the query engine.
Query cost is roughly (series matched) × (samples per series in range). A dashboard panel touching 50K series × 24h × 4 samples/min = 288M samples to decode — that is what kills query nodes, not the number of dashboards.
Pull vs Push#
| Model | Who | Strength | Weakness | Who Pays |
|---|---|---|---|---|
| Pull (scrape) | Prometheus, VictoriaMetrics agent | up == 0 is a free liveness signal; server controls load; targets are dumb | Needs service discovery; short-lived jobs vanish before scrape; NAT/firewall friction | Platform team owns discovery config |
| Push | InfluxDB, M3, StatsD, OTLP, Prometheus remote-write | Works for batch jobs, IoT, serverless, cross-network | Server cannot throttle clients; a buggy client can flood ingest; no free liveness | Ingest tier eats bursts; needs per-tenant rate limits |
Staff default: pull inside a cluster, push across boundaries. Prometheus scrapes locally and remote-writes (push) to the central store. That hybrid is what Mimir, Thanos Receive, and VictoriaMetrics are built for.
Data Modeling — "The Entire Game"#
In Cassandra the game is the partition key. In a TSDB the game is the label set, because the label set is the primary key, and every distinct combination is a new series with its own memory cost forever (well, until it goes stale).
Cardinality Is Multiplicative#
series(metric) = product over labels of (distinct values of that label)
http_request_duration_seconds_bucket
service : 40
endpoint : 25 (per service, avg)
method : 4
status : 6
le : 12 (histogram buckets)
pod : 30 (per service, avg)
= 40 x 25 x 4 x 6 x 12 x 30 = 8.64M series <- from ONE metric
Drop pod (aggregate at recording rule): 288K series
Add user_id (50M users): game over
The rule is simple and non-negotiable: every label must have a bounded, known set of values, and the product must fit the budget.
| Label kind | Example | Verdict | Where it goes instead |
|---|---|---|---|
| Bounded enum | method, status_class, region, tier | Label | — |
| Bounded infra | pod, instance, node | Label, but churns on deploy; aggregate away in rollups | Recording rules strip it |
| Semi-bounded | endpoint / route | Label only if templated (/users/:id), never raw paths | Normalize in middleware |
| Unbounded ID | user_id, order_id, trace_id, ip | Never a label | Logs, traces, exemplars, OLAP store |
| Free text | error_message, query | Never a label | Logs |
🎯 Staff Move: "I'd normalize the route to its template —
/orders/:id, not/orders/83712— in the HTTP middleware, because one raw path label is enough to OOM the cluster. And I'd attach the trace ID as an exemplar, not a label, so an engineer can jump from a p99 spike to one slow trace without paying for 50M series."
Metric Types and What They Cost#
| Type | Stores | Aggregatable across instances? | Gotcha |
|---|---|---|---|
| Counter | Monotonic total | Yes — sum(rate(x[5m])) | Resets on restart; always wrap in rate() |
| Gauge | Point-in-time value | Sum/avg/max are fine | Sampled — misses spikes between scrapes |
| Histogram | Counts per bucket (le) + sum + count | Yes — sum buckets, then quantile | Bucket count × other labels = cardinality; pick 8–15 buckets |
| Summary | Client-computed quantiles | No — cannot aggregate p99 across pods | Avoid for anything fleet-wide |
| Native/exponential histogram | Sparse exponential buckets | Yes | Prometheus 2.40+; far fewer series than classic histograms |
The single most common correctness bug in metrics design is averaging percentiles:
pod A: p99 = 50ms (10,000 req)
pod B: p99 = 900ms (100 req)
avg(p99) = 475ms <- meaningless
true fleet p99 ~= 55ms
Correct:
histogram_quantile(0.99,
sum by (le, service) (rate(http_request_duration_seconds_bucket[5m])))
Schema Design in InfluxDB and TimescaleDB#
InfluxDB (1.x/2.x) splits a point into tags (indexed, part of series identity) and fields (values, not indexed). The same cardinality rule applies to tags. The classic mistake is making device_id a field (then you can't filter on it efficiently) or making reading_id a tag (series explosion).
# line protocol: measurement,tags fields timestamp
cpu,host=web-17,region=us-east usage_user=12.4,usage_sys=3.1 1735689600000000000
TimescaleDB is PostgreSQL with hypertables: one logical table auto-partitioned into time-based chunks. Because it is relational, high-cardinality dimensions are just columns with a B-tree index — which is exactly why it wins the "analytics over time with joins" intent.
CREATE TABLE readings (
time TIMESTAMPTZ NOT NULL,
device_id BIGINT NOT NULL,
site_id INT NOT NULL,
temp_c DOUBLE PRECISION,
humidity DOUBLE PRECISION
);
SELECT create_hypertable('readings', 'time', chunk_time_interval => INTERVAL '1 day');
-- Columnar compression on chunks older than 7 days, segmented by device
ALTER TABLE readings SET (timescaledb.compress,
timescaledb.compress_segmentby = 'device_id',
timescaledb.compress_orderby = 'time DESC');
SELECT add_compression_policy('readings', INTERVAL '7 days');
-- Continuous aggregate: incrementally maintained hourly rollup
CREATE MATERIALIZED VIEW readings_1h WITH (timescaledb.continuous) AS
SELECT time_bucket('1 hour', time) AS bucket, device_id,
avg(temp_c) AS avg_temp, max(temp_c) AS max_temp, count(*) AS n
FROM readings GROUP BY bucket, device_id;
SELECT add_retention_policy('readings', INTERVAL '30 days');
Chunk sizing rule of thumb: size chunk_time_interval so the most recent chunk(s) across all hypertables fit in ~25% of RAM — the hot chunk's indexes must stay in memory or insert latency degrades.
Rollups — Store Aggregates, Not Averages#
A rollup must preserve the ability to re-aggregate. Store sum, count, min, max (and histogram buckets), never just avg — an average of averages is wrong the moment the underlying counts differ.
rollup_5m(series) = { sum, count, min, max, last } per 5-minute window
avg over any window = sum(sum) / sum(count) <- correct
avg(avg) <- wrong unless counts equal
The Tunable Tradeoff — Resolution × Retention × Cardinality#
Every TSDB design is a point in a three-dimensional budget. Pick two; the third is the bill.
storage_bytes/day ~= active_series x (86,400 / scrape_interval_s) x bytes_per_sample
1M series, 15s, 1.5 B/sample ~= 1M x 5,760 x 1.5 ~= 8.6 GB/day (raw)
10M series, 15s ~= 86 GB/day -> 15d raw ~= 1.3 TB
10M series, 60s ~= 21.6 GB/day
memory (head) ~= active_series x 3-8 KB (Prometheus; VictoriaMetrics ~1 KB)
10M series ~= 30-80 GB RAM on one Prometheus -> shard or move to Mimir/VM
| Lever | Turn it up | Who pays |
|---|---|---|
| Resolution (scrape 15s → 5s) | 3× samples, 3× ingest CPU, better spike detection | Platform budget; queries over long ranges get 3× slower |
| Retention (15d → 13 months raw) | ~26× storage; year-over-year comparisons at full fidelity | Finance; object storage makes this cheap-ish, index makes it slow |
Cardinality (add pod, customer_tier) | Multiplicative series growth, RAM, query fan-out | Every other tenant on the shared cluster |
| Replication (HA pair, RF=3 ingesters) | 2–3× ingest cost; survives node loss | Platform; dedup adds query complexity |
Tiered retention is the Staff default:
| Tier | Resolution | Retention | Storage (10M series) | Serves |
|---|---|---|---|---|
| Hot | 15s raw | 15 days | ~1.3 TB | Alerting, incident debugging |
| Warm | 5m rollup | 90 days | ~50 GB × 5 aggregates | Capacity trends, SLO reports |
| Cold | 1h rollup | 2 years | ~15 GB × 5 aggregates | YoY planning, finance |
🎯 Staff Move: "I'd ask product whether anybody has ever looked at 15-second resolution data older than two weeks. In every org I've seen, the answer is no — so raw retention is 15 days and everything beyond is 5-minute rollups. That single decision is usually a 10× storage reduction, and I'd have the SRE lead sign off on it because they're the ones who'll miss the data in a post-mortem."
Anti-Patterns — What Kills Time Series DB Deployments#
1. Unbounded Labels (The Cardinality Bomb)#
A developer adds customer_id to a latency histogram. Series count goes from 400K to 90M in the hour after deploy. Head memory OOMs, Prometheus restarts, WAL replay takes 15 minutes, and during those 15 minutes every alert in that cluster is silent. This is the #1 TSDB outage pattern, industry-wide.
Guardrail: per-tenant and per-metric series limits at ingest (Mimir max_global_series_per_user, VictoriaMetrics -maxLabelsPerTimeseries and series limiters, Prometheus sample_limit per scrape job). Reject, don't absorb.
2. Using the TSDB as an Event Store#
Writing one series per request, per order, or per click — "we'll query it later." A TSDB is optimized for many samples per series, not many series with one sample. Events belong in logs, Kafka, or a columnar OLAP store.
3. Averaging Percentiles and Storing Summaries#
Summaries compute quantiles client-side and cannot be aggregated. Fleet-wide p99 dashboards built on avg(summary_quantile) are silently wrong because every pod gets equal weight regardless of traffic: one idle, slow canary pod can drag the "fleet p99" up by 10×, and a few fast, busy pods can hide a real regression. The SLO report is fiction, and nobody notices until a customer does.
4. Short-Lived Series Churn#
Kubernetes pods, CI jobs, and autoscaled workers each create new pod= label values. A cluster with 10K pods redeploying twice a day creates 20K new series per metric per day. Steady-state count looks fine; the index over the retention window does not. Watch prometheus_tsdb_head_series_created_total rate.
5. Monitoring Inside the Blast Radius#
Prometheus running in the same Kubernetes cluster, on the same nodes, behind the same DNS, alerting through the same network path as the service it watches. When the cluster fails, the pager is silent. Classic correlated failure.
6. Unbounded Dashboard Queries#
A "last 30 days" dashboard with sum by (pod) over 500K series, auto-refreshing every 10s, opened by 40 engineers during an incident. Query nodes OOM exactly when they are needed most. Fix: query limits (max_samples, max_fetched_series_per_query), recording rules for heavy panels, and results caching.
7. Metering Revenue from a TSDB#
TSDBs deduplicate HA pairs, drop out-of-window samples, and downsample. Billing from sum(increase(api_calls_total[30d])) means under-billing on every counter reset edge case and no audit trail. Billing reads from an exactly-once event pipeline.
The Technology Landscape — Head-to-Head Comparison#
| Dimension | Prometheus | Thanos / Mimir / Cortex | VictoriaMetrics | M3 | InfluxDB 3 | TimescaleDB |
|---|---|---|---|---|---|---|
| Shape | Single node, local disk | Prometheus + object storage, horizontally scaled | Single binary or cluster (vminsert/vmstorage/vmselect) | Distributed M3DB + coordinator + aggregator | Columnar on Arrow/Parquet/object storage | PostgreSQL extension |
| Scale ceiling | ~5–10M active series per instance (RAM bound) | 100M–1B+ series | 100M+ series; low RAM per series | Billions of series (Uber-scale) | High ingest; cardinality far less painful | Tens of TB per node; multi-node needs care |
| Bytes/sample | ~1.3–2 | ~1.3–2 (+object store) | ~0.4–1 (claimed) | ~1.5 | Columnar Parquet | 90%+ reduction with native compression |
| High cardinality | Weak | Better, still index-bound | Better | Better | Strong (columnar, no series index) | Strong (it's just rows) |
| Query | PromQL | PromQL | MetricsQL (PromQL superset) | PromQL/Graphite | SQL, InfluxQL | Full SQL + joins |
| Ops burden | Low | High (many components) | Low–medium | High | Medium (or managed) | Medium (Postgres skills) |
| Pick when | Per-cluster scrape + alerting | Org-wide long-term metrics, multi-tenant | Cost-efficient Prometheus-compatible at scale | Very large, aggregation-heavy, already invested | IoT/events with high cardinality, SQL users | Time series that must join business data |
Managed/vendor options — Datadog, Grafana Cloud, Amazon Managed Prometheus / Timestream, Google Cloud Monitoring — collapse ops to zero and typically price per active series, per custom metric, or per sample ingested. That pricing model is exactly aligned with the cardinality problem, which is why an unreviewed label change can show up as a five- or six-figure invoice line.
Patterns#
Pattern 1: Edge Prometheus + Central Long-Term Store (the Staff default)#
- Alerting stays local — a central-store outage or WAN partition does not silence pages.
- HA pairs scrape identically; the central store deduplicates by
cluster+__replica__labels. - Object storage makes 13-month retention cheap (~$0.02/GB-month).
- When to use: more than one cluster, or retention beyond ~30 days, or more than ~5M series total.
Pattern 2: Federation / Recording-Rule Rollup#
A global Prometheus scrapes only pre-aggregated recording-rule series (job:http_requests:rate5m) from per-cluster Prometheus instances. Cheap and simple; loses raw drill-down. Use for small orgs (<5 clusters) that need a global SLO view without running Mimir.
Pattern 3: Push via Kafka for IoT and Events#
Devices → MQTT/HTTP gateway → Kafka (partitioned by device_id) → consumers batch-write to InfluxDB/TimescaleDB/VictoriaMetrics. Kafka absorbs bursts, allows replay after a TSDB outage, and fans out to the OLAP store. Use when producers are untrusted, bursty, or out-of-order. See Apache Kafka.
Pattern 4: Streaming Pre-Aggregation#
Instead of storing per-pod series and aggregating at query time, aggregate in the pipeline (M3 Aggregator, Flink, or vmagent stream aggregation) and store only the rollup. Cuts series 10–100×; loses the ability to drill into a single pod after the fact. Use for high-churn fleets where per-instance history is never queried. See Flink & Stream Processing.
Pattern 5: Dual-Store — TSDB for Health, OLAP for Questions#
Aggregates go to the TSDB for dashboards and alerts. Raw events with full dimensions (customer, SKU, experiment arm) go to ClickHouse/Druid/BigQuery. The TSDB answers "is it broken?"; the OLAP store answers "for whom?" Use whenever someone asks for per-customer latency.
Scaling#
Where Each Engine Hits the Wall#
| Scale (active series) | Samples/sec @15s | Architecture | What breaks next |
|---|---|---|---|
| < 1M | < 67K | One Prometheus HA pair per cluster, 15–30d local retention | Nothing; don't over-build |
| 1–10M | 67K–670K | Functional sharding (Prometheus per team/domain) + Thanos sidecar or remote-write | Head RAM (30–80GB), WAL replay time, cross-shard queries |
| 10–100M | 0.7–6.7M | Mimir / Cortex / VictoriaMetrics cluster, ingesters RF=3, object storage | Query fan-out, compactor lag, per-tenant noisy neighbors |
| 100M–1B+ | 6.7–67M | Multiple cells (per region/business unit), streaming aggregation, strict quotas | Org governance — whose series are these, and who pays? |
Sharding Strategies#
- Functional sharding — one Prometheus per team or domain. Simple, clear ownership, but cross-team queries require a global layer.
- Hash sharding by series —
hashmodrelabeling splits targets across N Prometheus instances. Distributes load; every query fans out to N. - Ingester sharding in Mimir/Cortex — distributors hash each series to a consistent-hash ring of ingesters, RF=3. Shuffle sharding assigns each tenant a subset of ingesters so one tenant's cardinality bomb cannot touch every ingester.
Ingest and Query Numbers#
| Component | Rough capacity | Note |
|---|---|---|
| Prometheus single node | ~1M samples/sec ingest; 5–10M active series comfortably with 64–128GB RAM | RAM ~3–8KB/series |
| Mimir ingester | ~1.5–3M active series per ingester (typical sizing) | RF=3 triples cost |
| VictoriaMetrics vmstorage | Frequently reported ~1KB RAM/series | Lower RAM is its main selling point |
| Query on 10K series × 24h @15s | ~57M samples decoded, ~1–3s | Use recording rules below 1s |
| Remote-write overhead | ~1–2 bytes/sample on wire after snappy | Queue shards scale with lag |
Multi-Region#
Metrics are one of the few datasets where you should not replicate synchronously across regions. Each region runs its own ingest and alerting; a global query layer (Thanos Query, Mimir with federated querier, or Grafana with multiple data sources) fans out read-time. Cross-region replication is only for DR of long-term blocks, and object storage cross-region replication handles that asynchronously at ~$0.02/GB transfer.
🎯 Staff Insight: "If us-east is on fire, I want us-east's alerting to keep working even if the central store in us-west is unreachable, and I want us-west's alerting to not care at all. So alerting is region-local, storage is region-local, and only the dashboard query layer is global."
Failure Modes & Recovery#
1. Cardinality Explosion → Head OOM#
- Symptom: Prometheus/ingester RSS climbs vertically after a deploy; OOM-kill; restart loop.
- Root cause: A new label with unbounded values (
user_id, raw URL path,error_message). - Detection:
prometheus_tsdb_head_seriesrate-of-change > 20% in 10 min;cortex_ingester_memory_seriesper tenant;scrape_series_addedper job. - Fix: Drop the label with
metric_relabel_configs(labeldrop) or block the metric at the distributor; restart; accept the WAL replay gap. - Prevention: Per-scrape
sample_limit, per-tenant series limits, CI lint on instrumentation (reject labels matching*_id), a cardinality dashboard per team.
2. WAL Replay Blackout#
- Symptom: After a crash, Prometheus is "up" but serving no queries and evaluating no rules for 5–20 minutes.
- Root cause: Replaying a multi-GB WAL rebuilds the entire head index in memory.
- Detection:
prometheus_tsdb_wal_replay_duration_seconds; readiness probe failing; dead-man's-switch alert fires. - Fix: Wait, or fail over to the HA twin (which is why HA pairs exist).
- Prevention: Keep per-instance series under ~5M; enable memory-snapshot-on-shutdown (Prometheus 2.30+ feature flag) to cut replay; always run HA pairs.
3. Remote-Write Backlog#
- Symptom: Central dashboards show data 10+ minutes stale; edge Prometheus memory grows.
- Root cause: Central distributor throttling, network partition, or ingester slowness; remote-write queue shards max out.
- Detection:
prometheus_remote_storage_highest_timestamp_in_secondsminusprometheus_remote_storage_queue_highest_sent_timestamp_seconds> 120s;prometheus_remote_storage_samples_failed_total. - Fix: Scale distributors/ingesters; raise
max_shardson the sender; drop non-critical metrics temporarily. - Prevention: Alerting stays on the edge (so stale central data is an inconvenience, not an outage); WAL-based remote write retries for ~2h of backlog.
4. Query of Death#
- Symptom: Query nodes OOM or peg CPU; dashboards time out for everyone.
- Root cause: One query matching millions of series over weeks (
sum by (pod) (rate(x[30d]))), often auto-refreshed. - Detection:
prometheus_engine_query_duration_secondsp99;cortex_query_frontend_queue_length; per-user query-sample counts. - Fix: Kill the query; identify the dashboard via query logs.
- Prevention:
query.max-samples(Prometheus default 50M), per-tenantmax_fetched_series_per_query, query-frontend splitting by day with results cache, recording rules for expensive panels.
5. Silent Monitoring (Correlated Failure)#
- Symptom: Nothing. No alerts during an outage. Discovered from customer tickets.
- Root cause: Monitoring shares the failure domain — same cluster, DNS, network, or Alertmanager route to a paging provider that is itself down.
- Detection: Only external: a dead-man's-switch (
vector(1)alert that must always fire) received by an external service that pages if it stops. - Fix: Page from the external watchdog; move monitoring off the failing substrate.
- Prevention: Run monitoring on a separate cluster or account, cross-region meta-monitoring pairs (us-east watches us-west), quarterly game day that kills the monitoring stack.
Operational Reality Matrix#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Cardinality bomb | prometheus_tsdb_head_series +20%/10m | Every tenant on the shared ingesters | Drop label at distributor; shuffle sharding | Platform on-call; emitting team fixes instrumentation |
| WAL replay | readiness probe; dead-man's switch | One Prometheus (HA twin covers) | Failover to twin | Platform |
| Remote-write lag | sent vs highest timestamp delta > 120s | Central dashboards only | Scale ingest; shed low-priority metrics | Platform |
| Query of death | query p99, frontend queue length | All dashboard users | Kill query; per-tenant limits | Platform; dashboard owner fixes query |
| Silent monitoring | External watchdog missing heartbeat | Whole org is blind | External paging path | Platform + SRE leadership |
| Clock skew on push clients | out_of_bounds / out-of-order sample rejections | One producer's data | NTP enforcement; OOO window | Producer team |
When to Use vs. Alternatives#
| Need | Pick | Why |
|---|---|---|
| Kubernetes/service metrics + alerting, < 5M series | Prometheus HA pair per cluster | Zero dependencies, PromQL, ecosystem |
| Org-wide metrics, 13-month retention, multi-tenant | Mimir / Thanos / VictoriaMetrics cluster | Object storage economics, horizontal scale, tenant limits |
| Cost-sensitive at 50M+ series | VictoriaMetrics | Lower RAM/disk per series; simpler ops than Mimir |
| IoT telemetry with SQL users, per-device queries | TimescaleDB or InfluxDB 3 | Cardinality is just rows/columns; SQL |
| Time series joined with business entities | TimescaleDB | It's PostgreSQL — joins, foreign keys, one system |
| Per-customer / per-experiment analytics, ad-hoc slicing | ClickHouse / Druid / Pinot | Columnar OLAP handles 100M+ distinct dimension values |
| Traces and per-request debugging | Tempo / Jaeger + exemplars | Not a metrics problem |
| Billing / usage metering | Kafka + exactly-once aggregation → OLTP ledger | Money needs exactness and audit, not a TSDB |
| Zero ops, < ~1M series, budget exists | Managed (Datadog, Grafana Cloud, AMP) | Buy until the bill exceeds a platform team |
When NOT to Use a TSDB#
- The question has an ID in it. "Latency for customer 8812" → OLAP or traces.
- Values must be exact and auditable. Billing, financial ledgers, compliance counts.
- Data is mostly updated, not appended. Current inventory, user profile state → OLTP.
- Events are sparse and unique. One sample per series is the worst case for every TSDB.
- You need joins across entities at query time — and you picked Prometheus. Use TimescaleDB or pre-join in the pipeline.
Operational Concerns#
What the On-Call Actually Does#
- Watches cardinality, not disk. Per-tenant series count, series churn, top-10 metrics by series (
topk(10, count by (__name__)({__name__=~".+"}))— run it on a replica, it's expensive). - Keeps the dead-man's switch green. An always-firing alert, received by an external service, paging if it goes silent for 5 minutes.
- Manages limits. Most platform tickets are "my metrics are being dropped" → a tenant hit its series limit. The runbook is: show them the top offenders, then raise the limit only with a sign-off.
- Watches compactor lag. In Mimir/Thanos, a stuck compactor means long-range queries touch thousands of uncompacted 2h blocks and slow down 10×. Alert when the oldest uncompacted block is > 24h old.
- Handles upgrades as rolling restarts of HA pairs, one replica at a time, waiting for WAL replay to finish before the next.
Key Metrics & Alerts#
| Metric | Healthy | Alert |
|---|---|---|
prometheus_tsdb_head_series | Stable ±5%/day | +20% in 10m |
rate(prometheus_tsdb_head_series_created_total[1h]) | < 1% of head/hour | Sustained > 5%/hour (churn) |
prometheus_rule_evaluation_duration_seconds | < 50% of eval interval | > 80% (rules falling behind) |
prometheus_notifications_dropped_total | 0 | Any increase |
prometheus_tsdb_compactions_failed_total | 0 | Any increase |
| Remote-write lag | < 30s | > 120s |
up == 0 per critical job | 0 targets | Any for > 2m |
| Dead-man's switch heartbeat | Every 1m | Missing 5m (external) |
Governance That Actually Works#
- Shadow → warn → enforce for new limits: log would-be rejections for 2 weeks, warn owners, then enforce.
- Cardinality report per team, weekly, with trend. Most explosions are caught by the owning team once they can see it.
- Instrumentation library defaults — the shared HTTP middleware ships templated routes and a fixed histogram bucket set, so the right thing is the easy thing.
Interview Application — Staff-Level Plays#
Which Case Studies Use Time Series DBs#
| Case Study | How the TSDB Is Used | Key Pattern |
|---|---|---|
| Metrics & Monitoring | The core store — this is the question | Edge scrape + central long-term store, cardinality limits |
| Ad Click Aggregation / Stream Processing | Serving layer for windowed aggregates | Pre-aggregate in Flink, store rollups only |
| Rate Limiting | Observability of limiter decisions, shadow-mode analysis | limiter_decisions_total{decision, policy} — never per-key labels |
| Auto-Scaling & Capacity | Signal source for scaling decisions and forecasting | 5m rollups × 13 months for seasonality |
| Circuit Breakers | Error-rate and latency histograms driving breaker state | Histograms aggregated by dependency |
| Ride Hailing & Delivery | Supply/demand per geo cell over time | Bounded geo cells (H3 res 7) as labels, never driver IDs |
Every System Design Question Has a TSDB Moment#
- URL shortener: "Click counts per link go to an OLAP store, not Prometheus — link ID is unbounded. Prometheus gets
redirects_total{status}for health." - Chat / messaging: "I'd track delivery latency as a histogram by region and client platform — 5 × 4 × 12 buckets = 240 series — and alert on the p99 burn rate against the SLO."
- Payment processing: "Authorization success rate by PSP and card network is my primary alert; the ledger itself never touches the TSDB."
- Job scheduler: "Queue depth and schedule lag per queue are gauges; per-job metrics would be unbounded, so job-level detail goes to logs with the job ID."
What Interviewers Probe#
| After You Say... | They Will Ask... | What They're Evaluating |
|---|---|---|
| "Prometheus for metrics" | "How many series? What happens at 50M?" | Do you size by cardinality or hand-wave? |
| "We'll add a customer label" | "How many customers?" | Do you recognize the multiplicative bomb? |
| "p99 latency dashboard" | "How do you compute fleet-wide p99?" | Histograms vs averaged summaries |
| "Keep a year of data" | "At what resolution? What does it cost?" | Retention as a priced product decision |
| "HA Prometheus" | "What if the whole cluster goes down?" | Meta-monitoring, failure domain separation |
| "Push from devices" | "What about late and out-of-order data?" | OOO windows, Kafka buffering, clock trust |
Common Interview Mistakes#
| What Candidates Say | What Interviewers Hear | What Staff Engineers Say |
|---|---|---|
| "We'll store every request in the TSDB" | Doesn't know TSDB vs event store | "Aggregates in the TSDB, events in Kafka → OLAP." |
| "Average the p99 across pods" | Statistically wrong SLOs | "Sum histogram buckets, then histogram_quantile." |
| "Add user_id so we can debug" | Will OOM the cluster | "Exemplars link metrics to traces; IDs never become labels." |
| "Prometheus scales horizontally" | Confuses Prometheus with Mimir | "Prometheus is single-node; I shard functionally or remote-write to Mimir." |
| "Keep everything forever, storage is cheap" | No cost model | "Raw 15d, 5m × 90d, 1h × 2y — a 10–20× reduction with product sign-off." |
| "Monitoring runs in the same cluster" | Correlated failure blind spot | "The watchdog lives outside the blast radius." |
L5 vs L6 vs L7 Responses#
| Scenario | L5 Answer | L6 / Staff Answer | L7 / Principal Answer |
|---|---|---|---|
| "Design monitoring for 2,000 microservices" | Prometheus + Grafana + Alertmanager | Per-cluster Prometheus HA pairs, remote-write to Mimir with RF=3 ingesters, 13-month object storage, per-team series quotas, local alerting, external dead-man's switch | Same architecture, plus: observability is a platform with a published SLO, chargeback per million active series, a label policy RFC, and a 2-year plan to consolidate the 4 existing stacks |
| "Metrics bill tripled" | Reduce scrape interval | Top-N cardinality report, drop pod label via recording rules, tier retention, enforce limits | Model vendor vs self-host crossover; showback per team so cost has an owner; decide which teams get premium retention |
| "Per-customer latency for enterprise SLAs" | Add customer_id label | Emit events to Kafka → ClickHouse for per-customer percentiles; TSDB keeps fleet health | Contractual SLA reporting becomes an audited data product owned by a named team, not a dashboard |
| "Monitoring was down during the outage" | Add more replicas | Separate failure domain, cross-region meta-monitoring, dead-man's switch | Org-wide correlated-failure review: which shared dependencies (DNS, IdP, paging vendor) can blind every team at once; game days twice a year |
The Staff TSDB Checklist#
- State the series budget: "~2,000 services × ~500 series each ≈ 1M active series per region; 15s scrape ≈ 67K samples/sec."
- Name the label policy: "Labels are bounded enums. Routes are templated. IDs go to exemplars and logs."
- Choose the topology: "Prometheus HA pair per cluster for scrape and alerting; remote-write to a Mimir-class store for global queries and long retention."
- Tier retention with a number: "15 days raw, 90 days at 5m, 2 years at 1h — roughly 10× cheaper than raw-forever."
- Design the failure posture: "Alerting is local, the watchdog is external, and each tenant has series and query limits so one team can't blind the rest."
- Draw the boundary: "Per-customer analytics and billing do not go here."
🎯 Staff Insight: Don't use a TSDB for anything with an unbounded identifier in the question, and don't bill from it. The strongest signal in a metrics interview is the moment you refuse a label.
Evaluation Rubric#
| Dimension | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| Sizing | Disk in GB | Active series, churn, samples/sec, RAM per series | $/series/month and the crossover vs vendor pricing |
| Data model | Metric names and labels | Cardinality math, histograms, route templating | Org-wide instrumentation standard and shared libraries |
| Architecture | Prometheus + Grafana | Edge scrape, central store, HA dedup, tiered retention | Cell-based platform per region/BU, multi-year consolidation plan |
| Failure | Replicas | Meta-monitoring, limits, query guards | Observability as its own failure domain with error budget and game days |
| Ownership | "SRE owns it" | Platform owns infra; emitting teams own their cardinality | Chargeback, quota exceptions process, platform/product contract |
Strong hire signals
| Signal | What It Sounds Like |
|---|---|
| Sizes by series | "10M series at ~4KB each is 40GB of head memory — that's my scaling trigger." |
| Refuses bad labels | "customer_id is not a label; it's a column in ClickHouse." |
| Correct percentile math | "I aggregate buckets, then compute the quantile." |
| Failure-domain thinking | "Who alerts us when the alerting is down?" |
| Prices retention | "Raw-forever costs ~12× more than tiered; I'd get SRE sign-off on the tradeoff." |
Lean no-hire signals
| Signal | Why It Misses the Bar |
|---|---|
| Treats the TSDB as a general event store | Wrong tool; will fail at the first high-cardinality question |
| Averages percentiles | Every SLO it produces is wrong |
| No limits or quotas in a shared cluster | One team's bug becomes everyone's outage |
| Monitoring colocated with the monitored system | Blind exactly when needed |
Common false positives
- Knowing Gorilla XOR bit layouts ≠ designing a metrics platform. Compression is solved; cardinality is not.
- Fluent PromQL ≠ capacity judgment. A beautiful query over 5M series is still a query of death.
- "We used Datadog" ≠ understanding the cost model. Ask what their bill was and why.
The Principal Lens#
Why L7 Sees This Problem Differently#
A Staff engineer sees a TSDB as a system to size and harden. A Principal sees it as the org's shared nervous system and its most elastic cost line. Every team emits metrics; almost no team pays for them. Observability spend commonly grows faster than the infrastructure it observes, because cardinality grows with services × instances × features while nobody's budget owns it. The L7 question is not "Mimir or VictoriaMetrics?" but "How do we make the cost of a label visible to the engineer who adds it, without making instrumentation so painful that teams stop measuring?"
The Org-Level Fault Line#
One central observability platform vs. every team runs its own stack.
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Central platform, free to use | Consistent tooling, global queries, one on-call | Tragedy of the commons; cardinality grows unbounded; platform team becomes the cost police | Platform team's budget and sanity |
| Central platform with quotas + chargeback | Cost has an owner; limits are defensible | Friction; teams game quotas; needs a billing pipeline | Product teams' budgets (correctly) |
| Per-team stacks | Autonomy; blast radius is the team | 10 stacks, 10 half-staffed on-calls, no cross-service view during incidents | Incident responders; the org at every multi-team outage |
| Vendor (Datadog/Grafana Cloud) | Zero ops, fast start | Per-series pricing makes cost scale with the worst team's instrumentation | Finance, surprised quarterly |
The Principal default: central platform, per-team quotas with showback in year 1 and chargeback in year 2, and a paved-road instrumentation library that makes the cheap thing the default.
Cost Model#
Assumptions: 15s scrape, ~1.5 bytes/sample, RF=3 ingest, object storage $0.023/GB-month, compute at ~$0.04/vCPU-hr and ~$0.005/GB-RAM-hr, loaded engineer cost ~$250K/yr. Vendor figures assume typical list pricing per active series or custom metric — directional only.
| Scale | Active series | Self-host infra/month | Platform headcount | On-call load | Vendor estimate/month |
|---|---|---|---|---|---|
| Startup | 500K | ~$1.5K (2 Prometheus pairs + object store) | 0.25 FTE | ~1 page/month | ~$3–8K |
| Growth | 10M | ~$15–25K (Mimir/VM cluster, 13mo retention) | 2 FTE (~$40K/mo loaded) | Weekly tuning, ~4 pages/month | ~$60–150K |
| Enterprise | 200M | ~$150–300K (multi-cell, per region) | 6–8 FTE | Dedicated rotation | Often $1M+; negotiated |
The crossover where self-hosting beats vendor pricing typically lands between 1M and 5M active series — but only if the org can staff 2+ engineers who want to run it. Below that, buy.
The 3-Year Evolution Path#
One-Way Doors vs Two-Way Doors#
| Decision | Reversibility | Cost to Reverse |
|---|---|---|
| Query language (PromQL vs SQL vs vendor DSL) | One-way-ish | Thousands of dashboards and alerts rewritten; 6–12 months |
| Instrumentation SDK (vendor agent vs OpenTelemetry/Prometheus client) | One-way | Touch every service; the most expensive migration in observability |
| Label naming conventions | One-way once dashboards depend on them | Dual-emit for a quarter, then break everything at once |
| Backend store (Thanos ↔ Mimir ↔ VM) behind remote-write | Two-way | Dual-write for a retention window; weeks |
| Retention tiers | Two-way going shorter, one-way once data is deleted | Lost history is gone |
| Scrape interval | Two-way | Config change; alert thresholds need retuning |
🧭 Principal Move: "The door I'd guard is the instrumentation SDK, not the database. If every service emits through OpenTelemetry or the Prometheus client with our shared middleware, I can swap backends in a quarter. If 400 services link a vendor agent, the vendor owns our roadmap."
The Standard I'd Write#
RFC-OBS-004: Metrics Instrumentation and Cardinality Standard
Scope: All services emitting metrics to the shared observability platform.
MUST
1. Emit via the approved SDK (OpenTelemetry or Prometheus client) through the
platform middleware; no direct vendor agents.
2. Use only bounded labels. Labels named *_id, *_uuid, email, ip, path (raw),
or free-text are rejected at ingest.
3. Record latency as histograms using the platform bucket set; summaries are
not accepted for fleet-level SLOs.
4. Stay within the team's active-series quota (default 250K per team per region).
SHOULD
5. Pre-aggregate per-instance metrics via recording rules when per-pod history
is not needed beyond 15 days.
6. Attach trace exemplars to latency histograms.
Exceptions: Filed with the observability platform team; approved by the platform
lead and the requesting team's director; time-boxed to 90 days; costed in the
request.
Rollout: shadow (log rejections) 4 weeks -> warn (dashboard + ticket) 4 weeks ->
enforce per region, canary region first.
Success metrics: series growth <= infra growth + 10%/yr; zero cardinality-caused
platform outages per quarter; p99 dashboard query < 2s; 100% of teams with a
visible showback line.
What I'd Tell the VP#
"Our monitoring cost is growing roughly twice as fast as our infrastructure because no team sees what their metrics cost. I'm proposing one shared platform with per-team quotas and a monthly cost report, moving to chargeback next year. This caps the bill at about $X per month at our projected scale, versus a vendor quote that grows with every new label. It also fixes the reliability risk: last quarter one team's change took alerting down for everyone, and quotas make that impossible. The ask is two engineers and a policy that directors sign."
Principal Interview Signals#
| Signal | What It Sounds Like |
|---|---|
| Prices the platform | "At 10M series we're at ~$20K/month self-hosted plus two engineers — the vendor quote crosses that at around 3M series." |
| Guards the right door | "I'll standardize the SDK before I standardize the database." |
| Designs cost ownership | "Showback first, chargeback in year two — cost without an owner only grows." |
| Org-level failure posture | "Which shared dependencies can blind every team at once? DNS, SSO, and the paging vendor." |
| Knows when not to standardize | "IoT telemetry stays on TimescaleDB; forcing it into PromQL helps nobody." |
Staff answers that L7 interviewers find insufficient:
- "We'll set per-tenant limits" — without saying who sets the quota, how exceptions work, or who pays for exceeding it.
- "Migrate to Mimir" — without a dual-write window, a dashboard migration plan, or a reason the org will fund it.
- "Use OpenTelemetry" — without recognizing it as a multi-year, every-team migration that needs an executive sponsor.
In the Wild#
Facebook — Gorilla#
Facebook built Gorilla as an in-memory write-through cache in front of its HBase-backed monitoring store, holding the most recent ~26 hours of data in memory. The 2015 paper introduced delta-of-delta timestamps and XOR float compression, averaging ~1.37 bytes per sample, and reported a 73× improvement in query latency. The key design choice was accepting that recent data matters most: almost all queries hit the last day, so optimize ruthlessly for that window and let older data live elsewhere.
Staff insight: The workload shape — recent data is hot, old data is rollups — justifies tiered retention. Quote the 1.37 bytes/sample number and then pivot immediately to cardinality, because compression is the solved half of the problem.
Uber — M3#
Uber built and open-sourced M3 (M3DB, M3 Coordinator, M3 Aggregator, M3 Query) after outgrowing earlier Graphite- and Cassandra-based metrics storage. M3 accepts Prometheus remote-write, stores data in a distributed TSDB with its own compression, and performs streaming aggregation so that downsampled series are computed at ingest rather than at query time.
Staff insight: At very large scale, you aggregate in the write path, not the read path. Say "streaming pre-aggregation" and name what you lose: per-instance drill-down after the raw retention window.
Netflix — Atlas#
Netflix built Atlas as an in-memory dimensional time series system for operational insight across its fleet, prioritizing fast queries on recent data and aggressive rollups for older data. It is designed around the reality that at Netflix scale the cost of keeping every dimension at full resolution is prohibitive, so dimensionality is managed deliberately.
Staff insight: Three companies, same conclusion — keep recent data hot and in memory, roll up aggressively, and treat dimensionality as a budget. That convergence is the argument to make out loud.
Practice Drill#
Prompt: "Our Prometheus-based monitoring for 3,000 services across 6 Kubernetes clusters keeps OOM-ing. Leadership wants 13 months of retention for capacity planning. Redesign it."
Staff Answer
First I'd measure: prometheus_tsdb_head_series per instance and the top 20 metrics by series count. My bet is 60–80% of series come from a handful of metrics with pod or unbounded labels, and churn from deploys. Immediate fix: drop unbounded labels via metric_relabel_configs, set sample_limit per scrape job, and shard each cluster's Prometheus functionally so no instance exceeds ~5M series. Architecture: keep a Prometheus HA pair per cluster for scraping and local alerting (15-day retention), and remote-write to a Mimir or VictoriaMetrics cluster backed by object storage. Size it: 3,000 services × ~3K series ≈ 9M active series ≈ 600K samples/sec at 15s — roughly 6–9 ingesters at RF=3. Retention is tiered: raw 15 days, 5m rollups 90 days, 1h rollups 13 months — about 12× cheaper than raw for 13 months, and capacity planning never needs 15-second data from last March. Per-team series limits with a shadow → warn → enforce rollout. External dead-man's switch so we know when monitoring itself is down. Per-customer latency questions go to ClickHouse, not here.
Why this is L6:
- Sizes by active series and samples/sec, not disk.
- Keeps alerting at the edge so the central store is not a single point of failure for paging.
- Treats retention as a tiered, priced decision and names who signs off.
- Adds limits and meta-monitoring — protects the platform from its own tenants.
What L7 adds:
- A showback report per team from month one, and a chargeback decision point at month twelve.
- Standardizes on OpenTelemetry/Prometheus client via shared middleware so the backend becomes a two-way door.
- A cost comparison against a managed vendor at 9M series, with the crossover stated explicitly.
Quick Reference Card#
Sample raw size: 16 bytes (int64 ts + float64 value)
Gorilla compression: ~1.37 bytes/sample (paper); Prometheus ~1.3-2 B
Head memory: ~3-8 KB per active series (Prometheus); VM ~1 KB
Prometheus ceiling: ~5-10M active series per instance
Head block: 2h in memory, WAL in 128MB segments
Samples/sec: series / scrape_interval (1M @15s = 67K/s)
Storage/day: series x 5,760 x 1.5 B (1M @15s ~= 8.6 GB/day)
Retention tiers: raw 15s x 15d | 5m x 90d | 1h x 2y
Histograms: 8-15 buckets; aggregate buckets, never percentiles
Default query guard: 50M samples per query (Prometheus)
Remote-write lag alert: > 120s
Dead-man's switch: always-firing alert, external receiver, page on silence
RED FLAGS
- Label containing an ID, raw path, IP, or free text
- avg() over a quantile
- Monitoring in the same failure domain as the monitored system
- Billing computed from a TSDB
- "Prometheus scales horizontally"