System Design Problem

Design a Time-Series Database

Commonly Asked By:InfluxDataDatadogPrometheusAWS

Interview Setup

Interview Prompt

Design a time-series database ingesting 10M data points/sec across 10M unique series, supporting 10K aggregation queries/sec with automatic downsampling and retention policies.

Clarifying Questions (ask before designing)

QuestionWhy it matters
Write-heavy (metrics/IoT) or balanced read/write?
  • Write-heavy favors LSM append-only
  • balanced may need columnar for analytics queries.
Query patterns: recent data only or historical range scans?Drives retention tiers: raw 7 days, 5-min rollup 90 days, 1-hr rollup 2 years.
Cardinality limits: what happens with tag explosion?
  • 10M series is manageable
  • unbounded labels (user_id as tag) kills any TSDB.
Out-of-order writes expected (network delays, batch uploads)?
  • IoT devices buffer and burst
  • TSDB must handle late-arriving data within tolerance window.

Scope

In scope

  • Write-optimized storage (LSM/columnar)
  • Downsampling
  • Rollup aggregation
  • Retention policies
  • Tag-based indexing
  • Out-of-order write handling

Out of scope (state explicitly)

  • Detailed frontend/UI pixel implementation
  • Org structure, staffing, and hiring plan

Functional Requirements

Start by confirming write vs query workload with your interviewer. Ask about retention policies, downsampling tiers, and whether alerting integration is in scope.

  • Ingest time-stamped data points at high throughput (metrics, IoT, financial ticks)
  • Write format: (metric_name, tags, timestamp, value)
  • Query by metric name + tag filters + time range
  • Aggregation queries: avg, sum, min, max, percentile over time windows
  • Downsampling: automatically roll up high-resolution data for older data
  • Retention policies: auto-delete data older than configured retention
  • Tag-based indexing: efficiently filter by any combination of tags
  • Support rate, derivative, moving average, and other time-series functions
  • Alerting integration: trigger alerts based on threshold queries

Non-Functional Requirements

Your interviewer will care most about write throughput and compression efficiency. TSDBs are 95% writes: call out Gorilla compression and block-based storage before they ask how you hit 10M points/sec.

  • High Write Throughput: 10M+ data points/sec ingestion
  • Low Query Latency: < 100ms for recent data; < 5s for large historical queries
  • Storage Efficiency: 2 to 4 bytes per data point after compression
  • High Availability: 99.99% for writes
  • Horizontal Scalability: add nodes for more throughput and storage
  • Durability: no data loss for committed writes
  • Write-Optimized: 95% writes, 5% reads (typical monitoring workload)

Capacity Estimations

Run this math before you size compression. Data points per second and retention window tell you disk growth; downsampling tiers reduce long-term storage by 10 to 100x.

MetricCalculationValue
Data points / sec (ingestion)From Data points / day ÷ 86400 (+ peak factor in value)10M
Unique time series (cardinality)Given (assumption documented in value)10M
Avg tags per seriesGiven (typical workload assumption)5 to 10
Raw data point sizeGiven (assumption documented in value)16 bytes (8B timestamp + 8B value)
Compressed data pointGiven (assumption documented in value)2 to 4 bytes (Gorilla compression)
Raw ingestion bandwidth10M x 16B160 MB/sec
Compressed storage / dayGiven~3.5 TB/day
Query rateGiven (assumption documented in value)10K queries/sec

Architecture Diagram

Walk your interviewer through the write path left to right. Time-series database architecture based on the Prometheus TSDB model: WAL → in-memory HEAD block → immutable blocks on disk → compaction → downsampling to object store. Each stage is labeled, including why metadata is written last for durability.

Loading...

Component Deep Dives

In the room: frame as append-only time-ordered writes, not OLTP: justify why PostgreSQL fails at 10M points/sec ingestion.

Gorilla Compression (Facebook's Time-Series Compression)

Delta-of-delta timestamps and XOR value encoding are the key insight: consecutive points in a series are similar.

Key insight: consecutive data points in a time series are similar.

Timestamp compression (Delta-of-Delta): t0 stored full (64 bits), t1 stores delta (variable bits), t2 stores delta-of-delta (if regular interval → 1 bit). Most points: 1 bit per timestamp when scrape interval is regular.

Value compression (XOR): v0 stored full (64 bits), v1 stores XOR with v0 (only changed bits), v2 same as v1 → 1 bit. For slowly changing metrics: ~2 bits per value.

Result: 16 bytes raw → 1.37 bytes average (12x compression). This is why TSDBs can handle 10M points/sec on modest hardware.

Block-Based Storage (Prometheus TSDB Model)

Block-based storage with 2-hour HEAD blocks is worth explaining: immutable blocks enable parallel scan, easy deletion, and independent caching.

Time axis is divided into 2-hour blocks. HEAD block: in-memory, accepts writes. After 2h: HEAD becomes immutable block on disk → new HEAD created. Compaction merges adjacent blocks up to max block duration (e.g., 48h). Benefits: writes go to single in-memory block (fast), immutable blocks are easy to cache/replicate/backup, deletion is entire-block delete, each block has its own index for parallel scan.

Event Bus Design (Kafka)

Topic: metrics-raw
  Partitions: 256 (partition by hash(metric_name + sorted_tags))
  Retention: 24h (ingestion buffer; blocks persisted to S3 after flush)
  Replication factor: 3, min.insync.replicas: 2
  Producers: Write Gateway, Telegraf agents, Prometheus remote_write
  Consumers: Ingester nodes (replay on crash before WAL catch-up)

Topic: compaction-events
  Partitions: 32 (partition by shard_id)
  Retention: 7 days
  Producers: Block Compactor (HEAD flushed → immutable block)
  Consumers: Downsampler (15s → 1m → 1h tiers), Retention Manager

Topic: out-of-order-drops
  Low volume audit stream when points arrive outside allowed window
  Consumers: cardinality guard alerts, client SDK telemetry

Write path: HTTP ingest → local WAL → batch to metrics-raw → ingester assigns shard
  Hot path returns 204 once WAL fsync completes; Kafka decouples burst from storage nodes

API Design

Write API

POST /api/v1/write
Content-Type: application/x-protobuf

# Prometheus exposition format:
cpu_usage{host="web-01", dc="us-east", env="prod"} 72.5 1710410700
memory_used_bytes{host="web-01", dc="us-east"} 8589934592 1710410700
http_requests_total{method="GET", path="/api", status="200"} 15234 1710410700

Query API (PromQL-style)

GET /api/v1/query_range
  ?query=avg(cpu_usage{dc="us-east"}) by (host)
  &start=1710324300 &end=1710410700 &step=60

Response:
{
  "status": "success",
  "data": {
    "resultType": "matrix",
    "result": [
      {"metric": {"host": "web-01"},
       "values": [[1710324300, "72.5"], [1710324360, "73.1"], ...]}
    ]
  }
}

Admin APIs

POST /api/v1/admin/retention   → Configure retention policies
POST /api/v1/admin/downsample  → Configure downsampling rules
GET  /api/v1/admin/status      → Storage node health, disk usage
POST /api/v1/admin/compact     → Trigger manual compaction

Common Error Responses

400 Bad Request: invalid input, missing required fields, or malformed JSON payload
401 Unauthorized: missing or invalid authentication token or API key
403 Forbidden: authenticated caller lacks required permissions for this resource
404 Not Found: requested resource ID does not exist
409 Conflict: duplicate write or version conflict, retry with a unique idempotency key
422 Unprocessable Entity: syntactically valid request failed semantic business validation
429 Too Many Requests: rate limit quota exceeded, client should honor Retry-After header
500 Internal Error: unexpected server failure, retry safely with an idempotency key
503 Service Unavailable: downstream dependency is unavailable or overloaded, retry with exponential backoff

Data Model

On-Disk Block Format (Prometheus TSDB)

Each block is a directory (named by ULID timestamp) containing: meta.json (time range, sample counts), index (inverted index mapping labels to series IDs + series metadata), chunks/ (Gorilla-compressed data files, max 512 MB each), and tombstones (deletion markers). A 2-hour block for 100K series with 15-second intervals ≈ 48M samples ≈ 66 MB.

In-Memory (HEAD Block)

HashMap<SeriesID, HeadChunk>

HeadChunk:
  labels: {__name__="cpu_usage", host="web-01", dc="us-east"}
  samples: Gorilla-compressed buffer (appends only)
  min_time, max_time, num_samples

Inverted Index (in-memory):
  label_name=label_value → BitSet of series IDs
  Fast intersection: host="web-01" AND dc="us-east" → bitwise AND

WAL (Write-Ahead Log)

Durability guarantees follow the append-only log pattern discussed in Write-Ahead Log (WAL) and Data Durability.

WAL segments (sequential, append-only):
  segment-00001:
    [record_type=SERIES, series_id=1, labels={...}]
    [record_type=SAMPLES, series_id=1, t=1710410700, v=72.5]
    [record_type=SAMPLES, series_id=1, t=1710410715, v=73.1]
On restart: replay WAL → rebuild HEAD block
WAL truncated: after HEAD block is flushed to persistent block

Fault Tolerance

High Cardinality Problem

Tags with unbounded values (user_id, request_id, IP) can create 100M unique series → 100 GB just for index. Solutions: cardinality limits (reject > 100K unique tag values), pre-aggregation at ingestion, separate high-cardinality data into ClickHouse/Druid, tag value allowlisting.

Write Path Durability

Client → Write Gateway → WAL (fsync) → ACK client → batch to storage node → storage node WAL + in-memory HEAD. Replication options: A) Write to N replicas independently (simple, slight gaps). B) Kafka buffer → multiple consumers (no data loss). C) Storage-level replication with tiered storage like Blob Storage (S3-like Object Storage) as used in Cortex and Thanos.

Query Performance Optimization

Block pruning (skip non-overlapping blocks), series filtering via inverted index, chunk LRU cache, query splitting into parallel sub-queries, step alignment to storage resolution, results caching in Redis, materialized views. Benchmark: 1M series, 6h query, 1m step → 30 seconds unoptimized → 200ms optimized.

Out-of-Order Writes

Options: Reject out-of-order (simple, loses data), accept with limited window (Prometheus v2.39+: within 1h), full out-of-order support (InfluxDB, VictoriaMetrics: higher memory). Recommendation: configurable out-of-order window (e.g., 30 min).

Additional Considerations

Comparison: TSDB vs General-Purpose DB

AspectPostgreSQLTSDB
StorageRow-based → 100+ bytes/pointColumn/chunk → 2-4 bytes/point (50x less)
Write patternB-tree index → write amplificationAppend-only → no index maintenance during write
CompressionNo time-series patternsGorilla compression → 12x
Throughput10K writes/sec max10M writes/sec

When PostgreSQL IS appropriate: low cardinality + low volume (< 1K writes/sec), need joins with relational data. TimescaleDB extension bridges the gap.

Multi-Tenancy & Isolation

Strategies: Tenant ID as top-level label (simple but noisy neighbor), per-tenant ingester (better isolation, more overhead), per-tenant storage (best isolation, easy billing). Rate limiting per tenant: write rate, series cardinality, query rate and timeout.

Worked Query Latency Example

Query: avg(cpu_usage) WHERE host="api-1" AND region="us-east" [last 24h, step=1m]

  • Index lookup: label matchers → ~50 series IDs (~2ms)
  • Block pruning: 24h at 2h/block → 12 blocks x 50 series = 600 block reads
  • Cache hit (80%): ~120 disk reads x 5ms = 600ms; cache miss adds ~2s
  • Aggregation: 1,440 points/series x 50 series in-memory (~50ms)
  • Total p99: ~800ms hot cache, ~2.5s cold, highlighting why downsampling to 5m resolution is critical for dashboards spanning weeks.

Interview Walkthrough

  • 25-minute cut

    Skip arch50/arch75 depth unless staff.

    • Frame as append-only, not OLTP random writes (5 min)
    • Gorilla compression: delta-of-delta + XOR encoding (6 min)
    • Block storage: HEAD block → flush → compaction (5 min)
    • WAL durability + inverted label index (5 min)
    • High cardinality is the silent killer: cap labels (4 min)
  • Frame as fundamentally different from OLTP: append-only time-ordered writes, not random read/write on indexed rows: justify why PostgreSQL fails at 10M points/sec.
  • Explain Gorilla compression: delta-of-delta timestamps and XOR value encoding exploit temporal locality: ~12x compression vs raw storage.
  • Walk through block-based storage: in-memory HEAD block → flush immutable 2-hour blocks → background compaction merges adjacent blocks.
  • WAL for write durability; inverted index (label → series IDs) for query-time series filtering via bitwise AND on label matchers.
  • High cardinality is the silent killer: unbounded tags like user_id or request_id explode the index; enforce limits and route high-cardinality data to ClickHouse.
  • Query optimization: block pruning by time range, chunk LRU cache, parallel sub-queries, and step alignment to storage resolution.
  • Common pitfall: using a TSDB as a general-purpose database with joins: it excels at time-range scans, not relational queries across entities.

Engineering Trade-offs

Write Path Optimization: LSM-Tree vs B-Tree for Time-Series

Time-series stores trade query flexibility against ingestion throughput, contrasting wide events with a narrow metrics schema.

B-Tree (PostgreSQL, MySQL): Good for random reads, but random writes cause page splits (write amplification). Not optimized for sequential time-range writes. At 10M writes/sec: too slow.

LSM-Tree (Cassandra, InfluxDB, Prometheus): Write to in-memory MemTable (sorted by timestamp) → flush to immutable SSTable → background compaction merges and compresses. Sequential disk writes, perfect for ordered timestamps, high write throughput (in-memory buffer absorbs bursts). Gorilla compression achieves 1.37 bytes per data point. Read amplification is the trade-off (check multiple SSTables).

Cardinality Explosion: The TSDB Killer

Low cardinality: cpu_usage{host="web-01"} → 100 series for 100 hosts.

High cardinality (dangerous): http_requests{path="/api/users/12345"} → 500M unique paths. Cardinality bomb: adding user_id tag → 1M users x 100 metrics = 100M series → OOM crash.

Solutions: label value limits (reject > 10K unique values), label normalization (use cohort labels instead of user_id), hierarchical metrics (roll up to coarser granularity), Bloom filter on series existence (1 GB for 100M series, O(1) check). Rule of thumb: max 1M unique series per Prometheus instance.

Downsampling and Retention: The Storage Cost Equation

Tier 1 (0-7 days raw): Every data point. 24 TB total.

Tier 2 (7-90 days, 1-min aggregates): MAX, AVG, MIN, P95 per minute. 28 GB (1000x compression).

Tier 3 (90 days - 2 years, 1-hour aggregates): ~1 GB. Total: ~25 TB vs ~10 PB without downsampling (400x reduction). Implementation: streaming aggregation via pipelines such as Stream Processing Basics (Flink + Kafka), or TimescaleDB continuous aggregates.

💬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...