When an interviewer asks you to design a metrics and monitoring system, they’re probing how you think about data pipelines, reliability, and user‑facing features. The problem is broad enough to let you showcase trade‑off analysis, but concrete enough that you can sketch a working design in 30‑45 minutes.

1. Clarify the scope and functional requirements

Start by asking clarifying questions. A typical set of functional goals looks like this:

  • Metric ingestion – agents or SDKs push numeric data points (e.g., CPU usage, request latency) at regular intervals.
  • Tagging and metadata – each point can carry dimensions like service name, region, or environment.
  • Retention and roll‑up – recent data is kept at high resolution; older data is down‑sampled.
  • Alerting – users define thresholds or anomaly rules that trigger notifications.
  • Dashboarding – a UI lets users query and visualize time‑series data.
  • APIs – programmatic access for both ingestion (/write) and querying (/query).

Non‑functional requirements (NFRs) usually include:

  • High write throughput – the system must handle bursts from many agents.
  • Low query latency – dashboards should feel snappy.
  • Durability – no metric loss, even under node failures.
  • Scalability – the design should grow from a single team to an enterprise.
  • Observability of the system itself – you need metrics about ingestion latency, storage usage, etc.

If the interview allows, ask whether the product is internal (e.g., for a single data center) or SaaS‑style (multi‑tenant). That will affect isolation and security decisions.

2. Identify core entities and data model

A time‑series database (TSDB) typically revolves around three concepts:

EntityDescription
MetricA named series, e.g., cpu.utilization. The name is often hierarchical (service.instance.cpu).
TagKey‑value pairs that add dimensions (host=web‑01, env=prod).
Data pointA tuple of (timestamp, value, tags).

Because you don’t know exact traffic numbers, use abstract variables:

  • N = number of distinct metric names.
  • M = average number of tags per metric.
  • R = write rate per agent (points per second).
  • A = number of active agents.

The total write volume is roughly N × M × R × A. Keeping the model flexible (e.g., allowing arbitrary tag sets) helps future‑proof the design.

3. Sketch a high‑level architecture

A clean way to break the system into layers is:

+-------------------+      +-----------------+      +-------------------+
|   Ingestion API   | ---> |  Stream Queue   | ---> |   Processing      |
|  (/write)         |      | (e.g., Kafka)   |      |  (validation,    |
+-------------------+      +-----------------+      |   aggregation)   |
                                                          |
+-------------------+      +-----------------+      +-------------------+
|   Query API       | <--- |  Index Store    | <--- |   Storage Engine  |
|  (/query)         |      | (TSDB indexes)  |      | (e.g., LSM‑tree)  |
+-------------------+      +-----------------+      +-------------------+

Ingestion API – lightweight HTTP endpoint that validates payloads and forwards them to a durable log (Kafka, Pulsar). Using a log decouples producers from consumers and gives replay capability.

Processing – a set of workers read from the log, apply down‑sampling, compute aggregates (e.g., 1‑minute averages), and write to the storage engine. You can also generate alerts here.

Storage Engine – a purpose‑built TSDB that stores raw points and roll‑ups. Column‑oriented stores (e.g., Parquet on object storage) work well for immutable data, while LSM‑tree based engines (e.g., RocksDB) excel at high‑write workloads.

Index Store – maintains inverted indexes for fast tag‑based lookups. A simple approach is a key‑value store mapping tag=value → list of metric IDs.

Query API – parses user queries, fetches relevant shards, merges results, and returns JSON or Prometheus‑compatible format.

4. Deep dive: Ingestion pipeline challenges

4.1 Back‑pressure and reliability

Agents may send bursts when a new deployment rolls out. If the queue fills, you need back‑pressure signals. Implement a circuit‑breaker on the API that returns 429 Too Many Requests when the log lag exceeds a threshold. Agents can then retry with exponential back‑off.

4.2 Schema evolution

Metric names and tags evolve. Store the schema in a registry service that maps metric names to allowed tag keys. Workers validate against this registry, but the system should tolerate unknown tags (store them as raw strings) to avoid breaking ingestion.

4.3 Compression and deduplication

Raw points are often repetitive. Apply delta‑encoding on timestamps and Gorilla‑style compression on floating‑point values before persisting. This reduces storage and network usage without sacrificing query speed.

5. Deep dive: Storage and query trade‑offs

5.1 Consistency vs. latency

If you write to a replicated log, you can choose between:

  • At‑least‑once – simple, but downstream workers must deduplicate.
  • Exactly‑once – requires idempotent writes and coordination, adding latency.

Most monitoring systems accept at‑least‑once semantics because duplicate points are harmless after aggregation.

5.2 Sharding strategy

Two common dimensions for sharding:

  • Metric‑first – hash the metric name; good for write locality but can cause hot shards if a single metric spikes.
  • Time‑first – partition by time bucket (e.g., hourly); spreads load evenly but requires cross‑shard merges for queries spanning many buckets.

A hybrid approach (hash + time) often balances the two.

5.3 Query path

Typical queries are range scans over a set of metric IDs filtered by tags. Use the index store to resolve tag filters to a list of metric IDs, then parallelize scans across shards. Merge the results in memory, applying any down‑sampling the client requested.

6. Alerting and anomaly detection

Alert pipelines can be built as a side‑car to the processing layer. After each aggregation window, evaluate user‑defined rules:

  • Static thresholds – e.g., cpu.utilization > 0.9 for 5 min.
  • Dynamic rules – percentile‑based or statistical models (e.g., moving‑average deviations).

When a rule fires, push a notification to a dispatcher that integrates with email, Slack, or PagerDuty. Keeping alert evaluation close to the aggregation step reduces latency.

7. Common follow‑up questions

QuestionTypical angle
How do you handle multi‑tenant isolation?Discuss namespace prefixes, per‑tenant quota enforcement, and separate index partitions.
What if a user wants millisecond‑resolution for the last hour but daily resolution older than that?Explain tiered storage: hot tier (raw points) for recent data, cold tier (down‑sampled) for older data, with background compaction jobs.
How would you support ad‑hoc queries like “show the 95th percentile latency for service X over the past week”?Show use of pre‑computed quantiles during aggregation or on‑the‑fly calculation using sketch algorithms (e.g., t‑digest).
What if the ingestion rate spikes 10× for a short period?Talk about autoscaling the ingestion workers, leveraging the log’s partitioning, and back‑pressure to producers.
How do you make the system observable?Mention internal metrics (ingestion lag, queue depth, storage hit‑rate) and a self‑monitoring dashboard.

When answering, keep the narrative concise and tie each decision back to a requirement you clarified earlier. If you need to rehearse the flow, a tool like Call Assistant can listen to you practice the answer aloud, keep follow‑ups on topic, and help you ground the story in your own résumé.

8. Sample answer template (45‑90 seconds)

"Sure, let me walk through a design for a metrics and monitoring platform. First, the functional goal is to ingest high‑frequency numeric data from many agents, store it efficiently, and let users query and set alerts. I’d expose a lightweight /write endpoint that validates JSON payloads and pushes them to a durable log like Kafka. Workers consume the log, apply compression, and write to a time‑series store that keeps raw points for the last hour and down‑samples older data. Tag indexes live in a key‑value store so queries can quickly resolve service=auth, env=prod to metric IDs. The /query API then parallel‑scans the relevant shards, merges results, and returns a JSON series. For alerting, I evaluate user‑defined thresholds right after each aggregation window and push notifications to Slack or PagerDuty. Non‑functional concerns include scaling the ingestion layer horizontally, using at‑least‑once semantics for simplicity, and keeping query latency under a second by sharding on metric‑hash plus time bucket. Trade‑offs involve choosing between exact consistency and lower latency, and handling hot‑spot metrics with a hybrid sharding scheme. Finally, the system is observable via internal metrics on ingestion lag and storage hit‑rate, which we surface in a self‑monitoring dashboard."

How to practice this

  1. Sketch the architecture on paper – draw the layers, label the data flow, and annotate where you’d place a queue or index.
  2. Explain the design out loud – use a timer and aim for a 60‑second summary; record yourself if possible.
  3. Iterate with feedback – ask a peer to play the role of the interviewer, focusing on the follow‑up questions above, and refine your answers based on their probes.

FAQ

  • What’s the difference between a time‑series database and a regular relational DB for metrics? A TSDB is optimized for append‑only writes, high‑cardinality tags, and range queries over time, while a relational DB excels at complex joins but struggles with the write volume and compression needs of monitoring data.

  • Why use a log (Kafka) instead of direct writes to the storage engine? The log decouples producers from consumers, provides durability, and enables replay for debugging or scaling the processing layer without losing data.

  • How can you reduce storage costs for long‑term metrics? Implement tiered retention: keep raw points for a short hot window, then down‑sample (e.g., average, max) into coarser buckets, and optionally move older buckets to cheap object storage.

  • When is exactly‑once delivery worth the extra latency? In scenarios where duplicate points could cause incorrect alerts or expensive recomputation, such as financial transaction monitoring. For most system metrics, at‑least‑once with idempotent aggregation is sufficient.

Frequently asked questions

What’s the difference between a time‑series database and a regular relational DB for metrics?

A TSDB is built for high‑write, append‑only workloads, supports high‑cardinality tags, and provides fast range scans over time. Relational databases handle complex joins but struggle with the volume and compression patterns typical of monitoring data.

Why use a log like Kafka instead of writing directly to storage?

The log decouples producers from consumers, guarantees durability, and lets you replay data for debugging or scaling. It also smooths bursty traffic by buffering writes before they hit the storage engine.

How can storage costs be reduced for long‑term metrics?

Apply tiered retention: keep raw points for a short hot window, then down‑sample into coarser buckets (averages, max) and optionally move older buckets to cheap object storage such as cloud‑based blobs.

When is exactly‑once delivery necessary?

Exactly‑once is useful when duplicate points could cause incorrect alerts or expensive recomputation, like financial transaction monitoring. For most system metrics, at‑least‑once with idempotent aggregation is sufficient.

#system design#metrics#monitoring#time series#architecture#a metrics and monitoring system