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)
| Question | Why it matters |
|---|---|
| Write-heavy (metrics/IoT) or balanced read/write? |
|
| 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? |
|
| Out-of-order writes expected (network delays, batch uploads)? |
|
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.
| Metric | Calculation | Value |
|---|---|---|
| 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 series | Given (typical workload assumption) | 5 to 10 |
| Raw data point size | Given (assumption documented in value) | 16 bytes (8B timestamp + 8B value) |
| Compressed data point | Given (assumption documented in value) | 2 to 4 bytes (Gorilla compression) |
| Raw ingestion bandwidth | 10M x 16B | 160 MB/sec |
| Compressed storage / day | Given | ~3.5 TB/day |
| Query rate | Given (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.
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 1710410700Query 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 ANDWAL (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 blockFault 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
| Aspect | PostgreSQL | TSDB |
|---|---|---|
| Storage | Row-based → 100+ bytes/point | Column/chunk → 2-4 bytes/point (50x less) |
| Write pattern | B-tree index → write amplification | Append-only → no index maintenance during write |
| Compression | No time-series patterns | Gorilla compression → 12x |
| Throughput | 10K writes/sec max | 10M 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_idorrequest_idexplode 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
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.