1. What It Is
Metrics and IoT data are write-heavy, time-ordered, and rarely updated — a different shape than relational rows. We use time-series databases or columnar stores tuned for append and range queries.
What:
A storage and query engine optimized for time-ordered numeric samples — each point is a metric name, value, timestamp, and optional label set (tags).
Primary purpose:
Ingest high-volume measurements cheaply, aggregate them over time windows, and alert when trends cross thresholds — without treating metrics like relational rows.
Usually used for:
Infrastructure monitoring, application APM, IoT telemetry, business KPI dashboards, and capacity planning — not general user profile storage.
2. Core Mental Model
We think in three dimensions when sizing a metrics pipeline:
⏱ Time
Every sample is anchored to a timestamp; queries are range scans, not primary-key lookups.
🏷 Labels
Dimensions like service=api, region=us-east — cardinality of label combinations drives memory cost.
📈 Aggregation
Raw points are rarely queried forever; rollups (sum, avg, max) at 1m/5m/1h tiers shrink storage and query cost.
In the room
Don't put metrics in PostgreSQL at scale. Mention downsampling (raw → 1min → 1hr aggregates), retention tiers, and cardinality explosion from high-cardinality labels (user_id as a tag kills Prometheus).
3. Why It Matters in HLD
Time-series data has write-heavy append patterns and time-range queries — specialized stores beat general SQL. Three lenses:
Needed When:
Problems mention dashboards, SLIs/SLOs, metrics pipelines, IoT sensors, or "billions of events per day" with time-range queries.
Avoids:
Storing every CPU sample in PostgreSQL rows — INSERT storms and index bloat kill the database.
Optimizes For:
Write throughput and time-range aggregation — not ACID transactions or complex JOINs.
4. Architecture & Data Flow
Walk ingestion and query as interview steps. Step 1 — Ingest: agents push metrics with timestamp + tags. Step 2 — Buffer: batch writes to LSM or columnar chunks. Step 3 — Index tags: inverted index on labels for filter. Step 4 — Query: time-range scan + aggregation (sum, rate, histogram). Step 5 — Retention: downsample old data; state roll-up policy.
In the room
Warn about label cardinality — high-cardinality tags (user_id on every metric) explode storage. Use bounded tags and aggregate at ingest.
5. Key Characteristics
Tag cardinality, compression, and downsampling — characteristics we compare across TSDBs:
| Metric Type | Behavior | Typical Query |
|---|---|---|
| Counter | Monotonically increases (requests_total) | rate() over window → requests per second |
| Gauge | Up or down (queue_depth, CPU %) | Instant value or avg over window |
| Histogram | Distribution of values (latency buckets) | p50/p99 from cumulative bucket counts |
- Append-only writes — samples are inserted, rarely updated in place.
- Compression — delta encoding and Gorilla-style algorithms exploit sequential timestamps.
- Retention policies — raw 7–15 days, aggregated months/years at coarser resolution.
- Pull vs push — Prometheus pulls from /metrics endpoints; InfluxDB and M3 use push ingest (Telegraf, remote write). M3 (Uber) shards time series across nodes with embedded etcd — mention when interviewer asks about Prometheus at multi-million-series scale.
6. Strategic Tradeoffs
Optimized time-range queries trade general-purpose SQL flexibility — we state both:
| Benefit | Cost |
|---|---|
| Write-optimized ingestion — millions of samples/sec via append-only paths and batch compression |
|
| Cheap long retention — rollups and columnar cold storage keep years of trends affordable | Query flexibility limits — ad-hoc joins across high-cardinality labels explode memory |
7. Failure / Bottleneck Awareness
Cardinality explosion, hot shards on single metric, and query timeout — we name mitigations:
Problem: A developer adds user_id as a metric label on HTTP requests. One million users → one million active time series; TSDB memory exhausts.
Mitigation: Cap label dimensions in code review; use logs or traces for high-cardinality IDs; aggregate before export; drop or hash labels at collection agent.
Problem: Dashboard queries 90 days of 1-second CPU samples per pod — billions of points scanned per refresh.
Mitigation: Pre-aggregate to 1m/5m rollups; query hot tier for recent detail, warm tier for historical trends; cache dashboard query results.
Problem: IoT devices with wrong clocks write samples into the future or distant past, corrupting rollups.
Mitigation: Reject samples outside acceptable window; use ingestion-time bucketing; reconcile with NTP on edge gateways.
8. Common HLD Usage
Monitoring dashboards, IoT telemetry, and financial ticks use time-series stores:
| Problem | Usage |
|---|---|
| Service monitoring dashboard |
|
| IoT sensor fleet |
|
| Ad impression analytics |
|
| SLO error budget tracking |
|
9. Decision Signals
Reach for TSDB when writes are append-only time-stamped events and queries are range aggregates:
- Queries always include a time range (last 1h, last 30d).
- Data is immutable measurements, not mutable entity state.
- Aggregations (sum, rate, percentile) matter more than fetching one row by ID.
- Write volume exceeds what a relational DB index can sustain on timestamp columns.
Top-K over millions of label combinations (e.g., "hottest articles today") needs global sort/aggregate across series — most TSDBs are optimized for per-series range scans, not leaderboard queries. Use a columnar warehouse, Redis sorted sets, or a pre-aggregated counter store instead.
🚫 Prefer OLTP or columnar analytics when:- You need UPDATE/DELETE on individual business records.
- Queries join many entity types with arbitrary filters unrelated to time.
- Exact row-level audit trails matter — use event logs or OLTP, not downsampled metrics.
11. Deep Dive (Optional)
Compression: delta and Gorilla encoding
Adjacent samples in a time series are often similar. Delta encoding stores the difference from the previous value instead of the full float. Timestamps with fixed scrape intervals compress further with delta-of-delta (if every point is +15s apart, you store one constant). Facebook's Gorilla paper applies XOR compression to floats — similar values share leading zero bits — achieving roughly 1–2 bytes per sample in practice.
Rollup Tier Design
Production systems rarely keep raw 1-second resolution forever. A typical tier ladder:
Raw (1s resolution) → retain 7 days → SSD hot tier Rollup 1 (1min avg/max) → retain 90 days → SSD warm tier Rollup 2 (1h avg/max) → retain 2 years → object storage / columnar Rollup 3 (1d avg/max) → retain 5+ years → cold archive (Parquet)
Background compaction jobs merge raw blocks into coarser buckets. Queries auto-select the finest tier that satisfies the requested time range — dashboards for "last 6 months" hit hourly rollups, not raw samples.
Prometheus Pull Model vs Push Gateway
Prometheus scrapes HTTP /metrics endpoints on a fixed interval (typically 15s). This gives uniform sampling and service discovery integration. Short-lived batch jobs cannot be scraped — they push to a Pushgateway, which Prometheus then scrapes. Interview pitfall: pushing every sample from every pod creates a single point of failure; pull is preferred for long-running services.
Histograms and Percentile Estimation
Storing every request latency as a raw sample is expensive at billions of QPS. Histograms bucket values (le=0.1, le=0.5, le=1.0 seconds) and increment counters per bucket. Percentiles are estimated from cumulative bucket counts — not exact, but stable under load. Choose bucket boundaries to match SLO thresholds (e.g., 100ms, 300ms, 1s for API latency).
TSDB vs Columnar Warehouse vs Search Index
- TSDB (Prometheus, InfluxDB, M3): low-latency range queries on metrics; weak at ad-hoc SQL joins.
- ClickHouse / BigQuery: excellent for analytics over event tables with time partition; higher query latency, cheaper at petabyte scale.
- Elasticsearch: full-text log search with time filters; not a replacement for numeric metric rollups at Prometheus scale.
In interviews, metrics dashboards → TSDB; business analytics over click events → Kafka + columnar warehouse; grep-able logs → search index. Many production stacks use all three with clear boundaries.
Capacity Sketch
10k pods x 500 metrics x 1 sample/15s ≈ 333k samples/sec. At ~2 bytes compressed per sample (optimistic), that is ~600 MB/min raw — retention and cardinality multiply this quickly. Always estimate active series count (unique metric+label combinations), not just sample rate.
Review
How helpful was this walkthrough?
Click a star to rate. We actively use this feedback to refine and update our system design content.
Discussion
Share your thoughts, ask questions, or help others.