Problem & requirements
The system targets internal operational metrics for a large company: CPU and memory, request counts, error rates, disk usage and similar. Business analytics and log aggregation are separate systems. Scale is around 100 million metric time series across many server pools. Data is kept for one year, but not at full resolution: raw points for about a week, one-minute rollups for a month and one-hour rollups for the rest of the year.
Alerts should fire to multiple channels (email, phone, PagerDuty, webhooks). Non-functionally the system must scale with the fleet, keep query latency low for dashboards, and be more reliable than what it monitors: a monitoring outage during an incident blinds the people fixing it. Flexibility to add new metric types without schema changes is also important.
Back-of-the-envelope
Suppose 1,000 server pools with 100 machines each and about 1,000 series per machine (including per-endpoint and per-label variants): 1,000 x 100 x 1,000 = 100 million series. Scraped every 10 seconds, that is 100M / 10 = 10 million data points per second, every second, all day.
A raw point is a timestamp and a float, about 16 bytes. 10M x 16 B = 160 MB/s, and times 86,400 that is roughly 13.8 TB per day uncompressed. Good encoding gets points down to under 2 bytes, cutting that to under 2 TB per day. Keeping raw data only seven days and then downsampling keeps the one-year footprint manageable. These numbers show why compression and retention policy matter more than query optimisation.
Data model and storage choice
A data point has a metric name, a set of labels (key-value tags such as host:web-12 and region:us-west), a timestamp and a value. A time series is a metric name plus a unique label set, and most queries pick series by labels and aggregate over a time range. The write load is constant and heavy; reads are spiky, driven by dashboards and incident investigations, and favour recent data.
A general relational database can store this but performs poorly: it is not tuned for this volume of tiny inserts, rolling-window aggregations require awkward SQL, and indexing on arbitrary labels is expensive. General NoSQL stores like Cassandra can work but demand deep schema expertise. Purpose-built time-series databases such as InfluxDB and Prometheus keep an inverted index from labels to series, store values in compressed blocks per series, and handle retention and downsampling natively.
High-level design
The pipeline has six parts. Metrics sources are application servers, databases and queues. The collector gathers points from them. A transmission pipeline moves data reliably. The time-series database stores it. A query service reads it for dashboards and alerts. The alerting system evaluates rules and notifies people, and a visualization layer such as Grafana draws graphs.
Putting Kafka between collectors and the database is the key scaling move. Without it, a slow or unavailable database causes collectors to drop data. With it, data is durably buffered, collectors and storage scale independently, and multiple consumers can read the stream for storage, real-time aggregation or anomaly detection. Partition topics by metric name so a consumer can aggregate a metric locally, and further by labels if a single metric is very hot.
Deep dive: pull versus push collection
In a pull model, collectors periodically scrape an HTTP endpoint on each target, finding targets through service discovery such as etcd or ZooKeeper; Prometheus works this way. Many collectors share the targets via consistent hashing. In a push model, an agent on each host sends metrics to collectors behind a load balancer; Graphite and CloudWatch work this way. Neither is universally better.
Each side wins somewhere. Pull gives easy health detection, because a target that stops answering is clearly down, and makes debugging simple since anyone can curl the endpoint. Push handles short-lived batch jobs that may exit before a scrape, and works when firewalls or complex networks block inbound connections. Pull usually runs over TCP; push can use UDP for lower overhead at the cost of loss. Push also needs authentication to stop arbitrary hosts from injecting data.
- Short-lived jobs: push wins (pull needs a push gateway).
- Firewalls and complex networks: push wins.
- Liveness detection and easy debugging: pull wins.
- Data authenticity: pull wins, targets are known in advance.
Deep dive: aggregation, querying and storage optimisation
Aggregation can happen in three places. At the collection agent it is cheap but limited to simple counters. In the ingestion pipeline, a stream processor such as Flink pre-aggregates before writing, which dramatically shrinks storage but loses raw precision and complicates late data. On the query side, raw data is kept and aggregated when queried, which is flexible but slower. Most systems mix these.
A query service in front of the database adds a cache for repeated dashboard queries and hides the database from visualisation tools. Storage optimisations do the heavy lifting: delta-of-delta encoding stores timestamps as tiny differences because scrape intervals are regular, and XOR-based float encoding exploits slow-changing values, as in Facebook's Gorilla. Downsampling rolls old data into coarser intervals, and cold storage keeps rarely queried history cheaply.
Alerting and visualization
Alert rules live as config files, often YAML, defining a query, a threshold and a duration, and are loaded into a cache. An alert manager periodically runs rule queries through the query service; when a condition holds it creates alerts. It must deduplicate and merge alerts, because one failing database can trip hundreds of rules, and it applies access control and routing. Alert state (inactive, pending, firing, resolved) is stored durably so restarts do not lose it.
Fired alerts go into Kafka, and alert consumers deliver them to email, SMS, PagerDuty or webhooks with retries so notifications are delivered at least once. Visualization should rarely be built in-house; Grafana and similar tools integrate with popular time-series databases and give mature dashboards for free. Building versus buying is a real interview discussion point here.
for duration before it moves from pending to firing, which filters out brief spikes. Alert state is persisted so a restarted alert manager neither re-pages nor forgets an open incident.Trade-offs, failure handling & wrap-up
The core trade-offs are pre-aggregation versus raw flexibility, push versus pull, and how aggressively to downsample. The failure story centres on the requirement to be more available than the monitored systems: run collectors redundantly, let Kafka absorb database outages, replicate the time-series database, and monitor the monitoring system from an independent stack so a blind spot is noticed.
In an interview, show you understand why a dedicated time-series store and Kafka buffer exist, quantify the write rate, and explain how compression and retention keep a year of data affordable. Mention that many teams adopt off-the-shelf tools such as Prometheus, Thanos or a managed service rather than building every layer.
Key numbers
Key terms
- Time series
- A sequence of timestamped values identified by a metric name and a unique set of labels.
- Label (tag)
- A key-value pair attached to a metric that lets queries filter and group series.
- Pull model
- Collectors scrape metrics endpoints on targets at fixed intervals.
- Push model
- Agents on hosts send metrics to collectors on their own schedule.
- Delta-of-delta encoding
- Compressing timestamps by storing the change in the difference between successive points, often zero.
- Downsampling
- Replacing high-resolution old data with coarser aggregates to save space.
- Alert deduplication
- Merging many alerts with a common cause into one notification to avoid paging storms.
- Cardinality
- The number of distinct series, driven by label combinations, which governs index size and cost.
Common mistakes
- Writing collectors straight to the database with no buffer, so storage outages lose data.
- Using high-cardinality labels such as user ID, which explodes the series count.
- Storing a year of raw 10-second data instead of downsampling.
- Declaring pull or push universally better instead of naming where each wins.
- Alerting without deduplication, causing alert storms during incidents.
- Running monitoring on the same infrastructure it watches.
Further study
- Gorilla: A Fast, Scalable, In-Memory Time Series Database (Facebook, 2015)
- Prometheus
- InfluxDB
- Grafana
- Google SRE Book: Monitoring Distributed Systems