Interview Setup
Interview Prompt
Design a distributed metrics aggregation system like Datadog or Prometheus at scale. Host agents collect CPU, latency, and custom counters, pre-aggregate locally, and push to a central store supporting rollups and percentile queries.
Clarifying Questions (ask before designing)
| Question | Why it matters |
|---|---|
| What is the ingestion rate in metrics per second and the unique series count? | Handling 10M points per second with millions of active series directly determines cardinality enforcement rules and storage engine selection. |
| Should the ingestion model use push (StatsD or OpenTelemetry) or pull (Prometheus scraping)? | Push models naturally accommodate short lived, auto-scaled containers, whereas pull models simplify agent complexity but require centralized service discovery. |
| What are the required retention tiers and downsampling policies? | Retaining raw 10-second data for 30 days, 1-minute rollups for 6 months, 1-hour rollups for 2 years, and optional 1-day rollups for 5 years drives storage capacity and compaction job scheduling. |
| What are the accuracy requirements for percentile calculations? | Calculating true global percentiles across distributed fleets requires mergeable histogram buckets or sketches, because averaging per-host percentiles is mathematically invalid. |
Scope
In scope
- StatsD-style push model
- Pre-aggregation at host agents
- Rollup storage tiers
- Percentile computation
- Capacity estimation with shown math
Out of scope (state explicitly)
- Application instrumentation SDK design
- Full distributed tracing system (covered in Distributed Tracing)
- On-call paging and escalation policy (covered in On-Call Escalation System)
Functional Requirements
Distributed metrics aggregation systems model all incoming operational telemetry as multidimensional time series. Client applications and infrastructure hosts emit counters, gauges, and histograms, which the aggregation platform ingests, groups by arbitrary label sets, downsamples across time windows, and exposes through analytical query interfaces for dashboards and alerting rules.
In an interview, proactively highlight that unbounded labels such as user_id on every request will destabilize the time series database index before raw volume does.
- Ingest metrics: Receive high frequency time series telemetry from hundreds of thousands of distributed hosts.
- Multi-dimensional aggregation: Aggregate metrics across arbitrary label dimensions using sum, average, min, max, and percentiles (p50, p95, and p99).
- Automated downsampling: Progressively compact historical metrics across rollup tiers, reducing raw 10-second data to 1-minute, 1-hour, and 1-day resolutions.
- Flexible range queries: Evaluate expressive analytical queries, such as average CPU utilization across region=us-east over the last 6 hours at 1-minute granularity.
- Alerting integration: Stream continuous metric updates into alerting engines to evaluate thresholds with sub-minute latency.
- Low-latency dashboarding: Deliver high-concurrency, sub second query responses for Grafana style visualization dashboards.
Non-Functional Requirements
At millions of active series and sub second collection intervals, cardinality explosion represents the primary threat to cluster stability. Clearly define a strict label budget, such as a maximum of 30 labels per series and a complete ban on unbounded identifiers like user_id, before sizing storage.
- High Throughput: Ingest and process over 10M data points per second at peak load.
- Low Query Latency: Deliver 99th percentile query latencies under 500 ms for standard operational dashboards.
- Tiered Long Term Retention: Preserve raw 10-second data for 30 days, 1-minute rollups for 6 months, 1-hour rollups for 2 years, and optional 1-day rollups for 5 years.
- System Scalability: Support 500,000 reporting hosts generating about 100M active time series.
- High Availability: Maintain 99.99% ingestion availability because telemetry outages directly blind alerting systems.
- Resilient Cardinality Control: Isolate and throttle runaway label combinations while preserving accepted operational telemetry.
Capacity Estimations
The active time series count dictates in memory index sizing, making it essential to calculate cardinality boundaries before selecting between federated Prometheus instances and a centralized time series storage cluster.
| Metric | Calculation | Value |
|---|---|---|
| Hosts reporting | Given workload assumption | 500K |
| Metrics per host | Given workload assumption | 200 |
| Total active time series | 500K hosts x 200 metrics/host, assuming each host metric forms one distinct series | 100M |
| Data points / sec | 100M active series ÷ 10-second reporting interval | 10M |
| Data point size | 8-byte timestamp + 8-byte value, excluding labels and encoding overhead | 16 bytes logical scalar sample |
| Ingestion throughput | 10M x 16 bytes = 160 MB/s | 160 MB/s |
| Logical raw payload / day | 160 MB/s x 86400 ≈ 13.8 TB | 13.8 TB before labels, encoding, and storage overhead |
| With downsampling | ≈81 TB compressed scalar payload + replication, indexes, metadata, block overhead, and headroom | ~500 TB production capacity budget |
Architecture Diagram
Metrics flow from host agents or scrapers into ingestion distributors through the selected push or pull ingestion protocol. In the steady state, distributors route samples directly to clustered time series ingesters, whose active blocks are written to durable object storage. A regional Kafka buffer activates during storage backpressure or ingestion bursts. Edge pre aggregation and strict cardinality enforcement protect downstream storage engines from unconstrained index growth.
In an interview, distinguish this telemetry engine from a dashboard visualization product such as the Real-Time Dashboard, as this architecture focuses specifically on the underlying ingestion, indexing, and aggregation substrate.
Component Deep Dives
Push vs Pull Ingestion
The aggregation engine coordinates edge collectors, ingestion distributors, time series storage partitions, and automated compaction workers. Each component manages a dedicated stage in the telemetry lifecycle.
Telemetry architectures choose between pull based polling and push based publishing based on workload ephemerality, network topology, and discovery requirements:
| Model | How | Pros | Cons |
|---|---|---|---|
| Pull (Prometheus) | Central server scrapes discovered targets on a configured interval | Server controls collection pace and discovers targets through service discovery | May require sharding, federation, or multiple Prometheus servers to scale across very large fleets |
| Push (StatsD / OpenTelemetry) ⭐ | Agents push batched metrics to an ingestion gateway | Scales across millions of ephemeral containers without service discovery overhead | Requires gateway load balancing to absorb traffic spikes, and agents must maintain gateway routing |
Specialized Time Series Storage Engines
Standard relational databases struggle under massive append volume and wide time-range aggregations. Specialized time series engines use append oriented write paths and compressed time series blocks. Prometheus style engines commonly use delta-of-delta timestamp encoding and XOR floating point compression, while exact storage internals vary by engine:
Relational Database Limitations (such as PostgreSQL): At 10M inserts per second, write-ahead logs and B-tree indexes saturate disk I/O bandwidth. Evaluating "avg cpu for last 6 hours across 1,000 hosts" triggers full table scans lasting minutes. Specialized Time Series Engine (VictoriaMetrics, InfluxDB, Cortex, or Mimir): 1. High write throughput through append-oriented write paths and compressed time-series block layouts. 2. Prometheus-style object-store TSDBs organize data into time-bounded blocks, with exact block and chunk behavior varying by engine. 3. High-density compression can use delta-of-delta timestamps and XOR float encoding, with ~1.5 bytes per sample treated as an illustrative target rather than a universal guarantee. 4. Background compaction and automatic downsampling organize historical retention tiers.
Tiered Downsampling Strategy
Recent telemetry demands fine grained resolution for active troubleshooting, whereas long term trend analysis requires only coarse aggregate statistics. Scheduled compaction jobs downsample historical data into progressive rollup tiers:
Resolution Tiers: Raw (10-second intervals): Retained for 30 days (Full operational detail for active debugging) 1-minute rollups: Retained for 6 months (6x reduction in data volume) 1-hour rollups: Retained for 2 years (360x reduction in data volume) 1-day rollups: Retained for 5 years (8,640x reduction in data volume) Illustrative compressed scalar payload budget: ~81 TB across the stated raw and rollup retention tiers Raw 30d: 10M points/sec x 1.5 bytes x 30 days ≈ 39 TB 1-minute 6mo: (100M series / 60s) x 1.5 bytes x 180 days ≈ 39 TB 1-hour 2y: (100M series / 3600s) x 1.5 bytes x 730 days ≈ 2.4 TB 1-day 5y: (100M series / 86400s) x 1.5 bytes x 1,825 days ≈ 0.25 TB Production capacity budget: ~500 TB, using roughly 3x replication plus label indexes, metadata, block overhead, and reserved headroom on top of the ~81 TB compressed scalar payload budget Approximate comparison: retaining the same 10M points/sec at raw resolution for 5 years would exceed 25 PB before labels and storage overhead
Mitigating Cardinality Explosion
Cardinality explosion occurs when dynamic metadata values, such as customer identifiers or request IDs, are assigned to metric labels. Because every unique label combination generates an independent time series, runaway cardinality rapidly exhausts index memory and crashes ingestion workers:
Problem Mechanics: Each unique request_id creates an entirely new time series. 10,000 requests/sec * 86,400 sec/day = 864M potential unique time series per day in the worst case when every request_id is new. System Impact: Ingester out-of-memory crashes, unbounded index expansion, and query timeouts. Mitigation Strategies: 1. Label cardinality limits: Reject metrics exceeding 10,000 unique values per label key. 2. Ingestion relabeling: Drop or hash high-cardinality metadata fields before indexing. 3. Active series caps: Enforce a strict quota of 1M active series per tenant. 4. Proactive monitoring: Track top metrics by series creation velocity and alert operators.
Kafka Buffer for Spike Absorption
The primary ingest path routes samples from hosts or scrapers through distributors to in memory ingesters and durable object storage blocks when storage is healthy. During sudden traffic surges or storage backpressure, a regional Kafka buffer absorbs burst volume and decouples ingestion workers from storage pressure:
topic: metrics-samples
purpose: Regional buffer for ingestion spike absorption
partitions: 256
partition_key: hash(tenant_id + metric_name + sorted_labels)
routing_behavior: Maps each tenant-scoped time series consistently to the same ingester shard
retention_hours: 6
retention_policy: Transient spike buffer only (durable blocks persist in S3)
replication:
factor: 3
min_insync_replicas: 2
acks: all
producer:
source: Regional aggregator after validating tenant and label cardinality
schema:
tenant_id: string
metric_name: string
labels: map[string]string
value: float64
timestamp: int64
type: counter | gauge | histogram
consumer_groups:
- name: ingester-writers
responsibility: Consume and forward burst-buffered samples to Cortex or Mimir ingesters
- name: cardinality-monitor
responsibility: Detect sudden label explosion from ingestion admission telemetry and trigger automated operator alerts
background_jobs:
- name: rollup-builder
responsibility: Build 1-minute and coarser rollups from finalized ingester blocks for long-term retention
pipeline_topology:
steady_state: Agent or Scraper to Distributor to Ingester to durable object storage, bypassing Kafka when storage is healthy
burst_path: Distributor to Kafka when storage write queues saturate, then Kafka consumers forward samples to ingesters
burst_ack: A burst-path write is acknowledged after Kafka confirms the configured replication level, while ingester persistence follows asynchronously
burst_absorption: Kafka buffers up to 10x the normal ingestion volume for the configured burst window
dead_letter_queue: metrics-samples-dlq (poison messages and non-retriable failures)
lag_alert: alert when consumer lag exceeds 30 secondsAPI Design
Domain Types and API Signatures
The aggregation platform exposes standardized endpoints for metric ingestion, instant metric evaluation, and historical range queries. The write endpoint carries tenant context in authentication metadata rather than treating a client-supplied tenant field as authoritative.
The remote write and query protocols use structured domain types representing tenant scoped labeled series, scalar samples, and histogram samples. The TypeScript model is a clean domain representation rather than a byte-for-byte Protobuf schema. Tenant identity is supplied by authenticated request context rather than trusted from the metric payload:
// Core multidimensional time series domain interfaces
export interface MetricLabel {
name: string;
value: string;
}
export interface MetricSample {
timestamp: number; // Unix epoch timestamp in milliseconds for remote write payloads
value: number; // Floating-point telemetry measurement
}
export interface HistogramBucket {
upperBound: number;
count: number;
}
export interface HistogramSample {
timestamp: number;
count: number;
sum: number;
buckets: HistogramBucket[];
}
export interface TimeSeries {
labels: MetricLabel[];
samples?: MetricSample[];
histograms?: HistogramSample[];
}
export interface WriteMetricsRequest {
series: TimeSeries[];
}
export interface AuthenticatedRequestContext {
tenantId: string;
}
export interface InstantQueryRequest {
query: string; // e.g., avg(cpu_usage{region="us-east"})
time?: number; // Unix evaluation timestamp in seconds
}
export interface RangeQueryRequest {
query: string; // e.g., rate(http_requests_total[5m])
start: number; // Range start timestamp in seconds
end: number; // Range end timestamp in seconds
step: number; // Query resolution interval in seconds
}Write Metrics (Prometheus Remote Write Compatible)
Host collectors and regional aggregators push compressed Snappy encoded Protobuf payloads containing tenant scoped time series and timestamped samples. Histogram data can be carried using the compatible histogram representation supported by the selected remote write version. The example below shows the logical request shape. Tenant identity comes from the authenticated request context, and the actual wire body is binary Protobuf:
POST /api/v1/write HTTP/1.1
Host: metrics-gateway.internal:9090
X-Tenant-ID: tenant_123
Content-Type: application/x-protobuf
Content-Encoding: snappy
X-Prometheus-Remote-Write-Version: 0.1.0
User-Agent: metrics-gateway/1.0
// Logical representation. The real request body is Snappy-compressed Protobuf.
TimeSeries {
labels: [
{ name: "__name__", value: "cpu_usage" },
{ name: "host", value: "h1" },
{ name: "region", value: "us-east" },
{ name: "env", value: "prod" }
],
samples: [
{ timestamp: 1710320000000, value: 85.2 }
]
}Time Series Query Endpoints
The query gateway evaluates instant queries against the latest evaluation timestamp and range queries across historical retention windows:
GET /api/v1/query?query=avg(cpu_usage{region="us-east"})&time=1710320000 HTTP/1.1
Host: query-gateway.internal:9090
X-Tenant-ID: tenant_123
GET /api/v1/query_range?query=rate(http_requests_total[5m])&start=1710316400&end=1710320000&step=60 HTTP/1.1
Host: query-gateway.internal:9090
X-Tenant-ID: tenant_123Common Error Responses
The ingestion gateway and query engine return standardized error payloads when rate limits are exceeded, label limits fail validation, or queries time out:
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 504 Gateway Timeout: search index shard responded slowly, narrow query parameters or retry
Data Model
Time Series Storage Layout
The storage subsystem organizes time series records into fixed duration compressed blocks and maintains an inverted postings index for multidimensional label querying.
Each tenant scoped time series is identified by a deterministic hash of the tenant ID and canonically sorted label pairs and stored in a sequence of compressed two-hour logical blocks. Physical chunk sizes vary by storage engine. The hash is a routing identifier rather than a guaranteed unique primary key, so the index retains enough canonical label data to detect and resolve rare hash collisions:
Series Definition:
cpu_usage{host="h1", region="us-east", env="prod"}
Storage Representation:
Series ID: hash(tenant_id + canonical_sorted_labels) = 0xABCD1234
Logical Block: 2-hour immutable block containing up to 720 samples at 10-second resolution
Physical chunks inside the block may be smaller depending on the storage engine
Timestamps: Delta-of-delta integer compression (Gorilla-style encoding)
Values: XOR double-precision floating point compression
Compression Efficiency: Approximately 1.37 to 1.5 bytes per recorded sample as an illustrative target for well-compressing numeric seriesInverted Postings Index for Label Filtering
An inverted postings index maps each label key-value pair to a sorted list of matching series IDs, allowing multidimensional queries to execute as efficient postings list intersections:
Inverted Postings Lists:
Label "region=us-east": Series IDs [0xABCD1234, 0xDEF56789, 0x9876FEDC]
Label "host=h1": Series IDs [0xABCD1234, 0x11223344]
Query Execution for cpu_usage{region="us-east", host="h1"}:
Step 1: Intersect postings lists for "region=us-east" and "host=h1"
Step 2: Identify matching Series ID (0xABCD1234)
Step 3: Retrieve index chunks matching the requested time window
Step 4: Decompress sample blocks and execute mathematical aggregationsFault Tolerance
The metrics aggregation platform must maintain ingestion continuity and query accuracy across hardware failures, traffic spikes, and network partitions:
| Concern | Solution |
|---|---|
| Ingestion spike | A regional Kafka cluster buffers sudden metric bursts while autoscaling groups provision additional ingestion workers. |
| TSDB node failure | In a quorum based ingester deployment, active state is replicated across 3 instances, while finalized historical blocks persist in durable object storage. Kafka based ingest-storage architectures can instead make Kafka the durable write path and let ingesters catch up asynchronously. |
| Data loss | Samples admitted to the ingester path are appended sequentially to an on-disk write ahead log before the configured durability acknowledgment. Process restarts can then reconstruct in memory state, while replica or quorum protection covers host and local disk failure. |
| High cardinality | Strict admission controllers enforce label quotas and reject any metric stream exceeding 100,000 unique active series per metric name. |
| Query overload | Query gateways enforce concurrency limits, apply a strict 30-second query timeout, and cache repeated dashboard queries. |
| Clock skew | The ingestion gateway accepts sample timestamps within a 5-minute window of server time, discards future-dated or severely lagged samples, and deduplicates retries using the tenant scoped series identity and sample timestamp. |
Additional Considerations
Alerting and Recording Rules
Alert evaluation runs as a separate horizontally scalable service that periodically evaluates rules against recent ingester data and historical query results. Recording rules pre-compute expensive expressions such as fleet-wide error rates or p99 latency every minute, reducing repeated dashboard and alert query cost. Rule evaluation is sharded by tenant and rule identity, with leader election or deterministic ownership preventing duplicate evaluations. Alert state is persisted or checkpointed so evaluator restarts do not repeatedly fire the same alert, while an external notification component handles routing, silences, and escalation.
Percentile Calculation at Scale
Production observability architectures require specialized mechanisms for accurate percentile calculation, crash recovery, multi tenant isolation, and trace correlation.
Calculating global percentiles across a distributed fleet requires mathematically mergeable data structures, because averaging per-node percentiles introduces severe mathematical error:
Mathematical Constraint: Averaging per-node p99 values produces false percentiles and masks latency spikes. Mergeable Aggregation Solutions: 1. Histogram Buckets (Recommended): Pre-define static latency boundaries (such as le=0.1, le=0.5, le=1.0) and sum counts across reporting nodes. 2. T-Digest: Streaming quantile algorithm that clusters sample distributions with bounded memory and approximately 1% error. 3. DDSketch: Fully mergeable sketch offering relative error guarantees across arbitrary percentile evaluations with bounded memory for a configured accuracy and value range.
Write Ahead Log (WAL) for Ingester Crash Recovery
To provide crash recovery without turning every query into a disk operation, ingesters record admitted samples to an append-only write ahead log on local NVMe before the configured durability acknowledgment:
Durability Lifecycle: 1. Sample arrives at the ingester and appends sequentially to the local disk WAL before the configured durability acknowledgment. 2. Sample updates active in-memory chunk block for real-time querying. 3. Every 2 hours, completed chunks seal, compress, and flush to persistent S3 object storage blocks. 4. On sudden process crash, a replacement ingester replays the 2-hour WAL in approximately 30 seconds. Process crash recovery can avoid sample loss for samples durably acknowledged to the WAL, while replica or quorum protection covers local disk or host failure. The stated replay time is an illustrative capacity target rather than a universal guarantee.
Multi Tenant Isolation and Rate Limiting
Multi-tenant observability clusters enforce strict ingestion quotas and query boundaries to prevent noisy tenants from degrading system health:
Per-Tenant Enforced Quotas: Maximum ingestion throughput: 100,000 samples per second Maximum active series count: 1,000,000 unique series Maximum label keys per series: 30 labels Maximum label value length: 2,048 characters Isolation Mechanisms: Tenant-aware consistent hash sharding, dedicated query worker pools, and billing usage telemetry.
Histogram vs Summary Semantics
Understanding the distinction between Prometheus style histograms and summaries is vital for designing reliable service level indicators. Histograms expose bucket counters that merge correctly across hosts, making them the standard choice for fleet-wide SLO percentiles. In contrast, client summaries pre-compute quantiles locally on each instance, making fleet aggregation statistically invalid. Modern platforms also attach exemplars containing distributed trace identifiers to histogram samples, allowing operators viewing the Real-Time Dashboard to jump directly from an anomalous latency bucket to the exact trace in Distributed Tracing.
Related Problems and Core Concepts
Deepen your understanding of observability, high throughput streaming, and time series storage through these related architectural designs:
- Real-Time Dashboard: High-performance visualization product consuming aggregated metrics for low latency dashboard querying and live alerting.
- Distributed Tracing: Correlates latency anomalies identified in metrics dashboards with end to end distributed span graphs across microservices.
- On-Call Escalation System: Routes actionable alert notifications triggered by metric threshold evaluations to on call engineering schedules.
- Time Series and Metrics Storage: Deep dive into columnar block layouts, Gorilla delta-of-delta compression, and inverted index postings lists.
- Stream Processing Basics: Windowing models, tumbling aggregations, and edge reduction patterns for massive telemetry pipelines.
- Observability and Distributed Tracing: Comprehensive foundations spanning metric collection, distributed tracing propagation, and structured logging.
- Back-of-the-Envelope Estimation: High-level capacity planning frameworks for quantifying network bandwidth, active series cardinality, and disk retention tiers.
Interview Walkthrough
When presenting this design during an interview, structure your pacing across core requirements, storage mechanics, cardinality mitigation, and mergeable percentiles:
- 25-minute interview pacing guide
Focus on core data structures and edge pre aggregation before diving into distributed storage details.
- Define the time series data model and establish strict cardinality limits (5 min)
- Compare pull based scraping against push based agent ingestion (6 min)
- Design ingester nodes with write ahead logging and object storage compaction (5 min)
- Formulate cardinality admission controls and tenant quotas (5 min)
- Explain mergeable histogram buckets and T-Digest percentile calculations (4 min)
- Start with the foundational data model, establishing that every time series consists of a metric name, a label set, and a sequence of timestamped values. Highlight immediately that label cardinality is the primary scaling bottleneck.
- Compare pull based scraping used by Prometheus against push based agents like StatsD or OpenTelemetry, defending push for ephemeral container fleets and pull for centralized target discovery.
- Route admitted samples through ingester nodes equipped with local write ahead logs so process restarts can recover durably acknowledged samples. Replica or quorum protection covers host failure.
- Enforce strict per-tenant label quotas, capping labels at 30 keys per metric and rejecting unbounded identifiers like user IDs to prevent runaway index growth.
- Calculate global percentiles using fixed histogram buckets or streaming T-Digests because averaging per-node percentiles is mathematically invalid.
- Shard cluster state by tenant or metric name hash with isolated query worker pools so one noisy tenant cannot starve platform resources.
- Quantify ingestion math: 10M samples per second at 16 bytes yields 160 MB/s of logical scalar sample payload before labels and encoding overhead. Compression, aggregation, and retention tiering substantially reduce the long term footprint.
- Highlight the common operational pitfall of applying identical time-to-live values to metrics, and demonstrate how jittered downsampling prevents database overload.
Engineering Trade-offs
Time Series Storage Engine Comparison
Architecting an enterprise metrics platform involves balancing write throughput against data freshness, local pull scraping against gateway push pipelines, and single node simplicity against horizontally scalable object backed clusters.
Selecting a time series storage backend requires evaluating clustering support, ingestion protocols, long term storage mechanisms, and query dialect compatibility:
| Feature | Prometheus | VictoriaMetrics ⭐ | InfluxDB | Mimir |
|---|---|---|---|---|
| Scalability | Single node | Clustered | Clustered or managed, edition dependent | Horizontally scalable |
| Ingestion | Pull only | Push and pull | Push | Push (remote write) |
| Storage | Local disk | Local and S3 | Local or object storage, edition dependent | Object store (S3) |
| Query language | PromQL | MetricsQL | InfluxQL / Flux | PromQL |
Ingestion Paradigm: Pull Based Polling vs Push Based Streaming
The choice between pull based and push based ingestion fundamentally shapes network topology, agent resource footprints, and service discovery mechanisms:
Pull-Based Ingestion (Prometheus Scrape Loops):
Mechanics: Central scraper periodically polls HTTP /metrics endpoints across registered targets.
Advantages:
- Server dictates collection pacing, which bounds the rate at which targets are scraped.
- Scrape failures provide a direct signal that a target is unreachable or unhealthy.
- Central scheduling makes unexpected client-side sending bursts less likely to overload the collector.
Disadvantages:
- Requires centralized service discovery such as the Kubernetes API or Consul to track dynamic endpoints.
- Poor fit for short-lived batch jobs and serverless functions that terminate between scrape intervals.
- Scaling requires complex hierarchical federation trees.
Push-Based Ingestion (StatsD and OpenTelemetry Gateways):
Mechanics: Lightweight client agents batch measurements and push asynchronously to ingestion gateways.
Advantages:
- Seamlessly supports ephemeral containers, mobile clients, and short-lived lambda workers.
- Agents pre-aggregate high-frequency telemetry at the edge, reducing network bandwidth by 10x.
- Simplifies network routing when clients reside behind private corporate firewalls.
Disadvantages:
- Ingestion gateways require dedicated buffering tiers (such as Kafka) to absorb traffic surges.
- Detecting node failure requires heartbeat leases rather than simple scrape timeouts.Percentile Evaluation: Histograms vs Summaries vs Streaming Sketches
Evaluating percentiles across distributed fleets involves significant trade-offs between mathematical precision, memory consumption, and query computational overhead:
Fixed Histogram Buckets (Prometheus):
Mechanics: Increments cumulative counter buckets partitioned by pre-defined boundaries.
Trade-offs:
✓ Buckets sum cleanly across thousands of nodes to calculate fleet-wide quantiles.
✓ Constant memory overhead per time series regardless of sample volume.
✗ Requires estimating latency distributions upfront, because bucket boundaries cannot be adjusted retroactively without data migration.
Client-Side Summaries:
Mechanics: Computes quantile estimates on the client instance over a sliding window.
Trade-offs:
✓ Highly accurate on a single host without pre-configured bucket boundaries.
✗ Summaries cannot be merged across hosts, which means averaging per-node p99 values produces mathematically invalid fleet percentiles.
Streaming Sketches (T-Digest and DDSketch):
Mechanics: Compress sample distributions into compact structures that can be merged across reporting nodes.
Trade-offs:
✓ Support arbitrary quantiles (p50, p90, p99, p99.9) with mergeable properties.
✓ Bounded memory, with observed error depending on the selected sketch and configuration.
✗ Higher CPU overhead during query evaluation compared to simple bucket additions.
HdrHistogram:
✓ Useful when a fixed latency range and high-resolution histogram are appropriate.
✗ Less flexible than T-Digest when distributions or percentile ranges vary widely.Storage Topology: Hierarchical Federation vs Object-Store TSDB
Scaling time series storage beyond a single machine requires selecting between hierarchical federated instances and centralized object-storage architectures:
Hierarchical Federation (Prometheus Scrape Trees):
Mechanics: Regional Prometheus instances scrape local clusters, while a global master scrapes rollups.
Trade-offs:
✓ Complete operational autonomy per region with zero shared storage dependencies.
✗ Cross-region queries must evaluate multiple hierarchical tiers, and pre-aggregated rollups discard fine-grained raw sample details.
✗ High administrative overhead maintaining hundreds of individual server configurations.
Object-Store Metrics Architecture (Cortex, Mimir, Thanos):
Mechanics: Cortex and Mimir provide distributed ingestion and query paths backed by object storage, while Thanos extends Prometheus with object-store retention and federated querying. Exact replication, Kafka usage, indexing, and read-path behavior depend on the selected engine and architecture mode.
Trade-offs:
✓ Separation of compute and storage allows horizontal scaling of query and ingest workers independently.
✓ Object storage (S3) provides nearly infinite, highly cost-effective long-term retention.
✗ Operational complexity requires managing multiple microservices (distributor, ingester, querier, compactor).
✗ Querying cold historical blocks introduces higher initial latency until indices are cached.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.