Core Concept

Time-Series & Metrics Storage

Purpose-built storage for timestamped measurements — counters, gauges, and histograms — with append-only ingestion, label cardinality discipline, rollup tiers, and retention policies tuned for observability and IoT workloads.


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.

Loading...

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 TypeBehaviorTypical Query
CounterMonotonically increases (requests_total)rate() over window → requests per second
GaugeUp or down (queue_depth, CPU %)Instant value or avg over window
HistogramDistribution 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:

BenefitCost
Write-optimized ingestion — millions of samples/sec via append-only paths and batch compression
  • Poor fit for OLTP — no arbitrary UPDATE/DELETE
  • point lookups by primary key are awkward
Cheap long retention — rollups and columnar cold storage keep years of trends affordableQuery 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:

🏷 Label Cardinality Explosion

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.

📊 Querying Raw at Scale

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.

⏰ Clock Skew & Late Data

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:

ProblemUsage
Service monitoring dashboard
  • Prometheus scrape + Grafana
  • 15-day hot retention, remote write to long-term store
IoT sensor fleet
  • Per-device timestamped readings
  • downsample raw 1s → 1min aggregates after 48 hours
Ad impression analytics
  • Counter + label dimensions (campaign, region)
  • cardinality caps on user_id labels
SLO error budget tracking
  • Histogram of request latency
  • burn-rate alerts on 5m and 1h windows

9. Decision Signals

Reach for TSDB when writes are append-only time-stamped events and queries are range aggregates:

🎯 Reach for time-series storage when:
  • 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.
Avoid TSDB when the problem is cross-series ranking:

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

Help Us Improve

How helpful was this walkthrough?

Click a star to rate. We actively use this feedback to refine and update our system design content.

Placeholder
Optional but highly appreciated!

Discussion

Share your thoughts, ask questions, or help others.

Loading comments...