Design a Distributed Metrics Monitoring & Alerting System
1. Problem Statement & Scope Clarification
System Mission
Design a planetary-scale, multi-tenant distributed Metrics Monitoring, Time-Series Storage, and Alerting Platform (equivalent to Prometheus/Thanos, Datadog, or Amazon CloudWatch) capable of ingesting tens of millions of metric points per second, executing real-time alert rules with sub-minute evaluation cycles, and delivering sub-second analytical dashboard queries over multi-year time horizons.
Functional Requirements
- High-Throughput Metric Ingestion (
PushMetric/BatchIngest): Continuously ingest structured time-series data points comprising metric name, 64-bit millisecond timestamp, 64-bit IEEE 754 float value, and multi-dimensional string key-value tag labels (e.g.,cpu.util{env="prod", service="payment", region="us-east-1"}). - Multi-Dimensional Tag Filtering & Analytical Queries: Support flexible label slicing and time-series aggregation functions (
avg,sum,rate,p95,p99) executed over custom sliding time windows. - Real-Time Sliding-Window Alert Evaluation: Continuously evaluate user-defined alert conditions (e.g.,
rate(http_requests_5xx[5m]) > 0.05) every 15–30 seconds, dispatching deduplicated notifications to Amazon SNS, PagerDuty, and Slack. - Automated Downsampling & Tiered Retention: Automatically roll up 1-second raw metrics into 1-minute and 1-hour summaries to optimize storage costs over a 5-year retention horizon.
- High-Cardinality Defense: Detect and strip unbounded dynamic tag dimensions (e.g.,
user_idororder_idinjected into metric tags) before they exhaust inverted index memory.
Non-Functional Requirements (SLAs & SLOs)
- Ingestion Throughput: sustained average (), with peak ingestion headroom of ().
- Query Latency (P95):
- Dashboard Queries on Recent Data (): .
- Range Queries on Historical Data (): .
- High Availability: ("five nines") uptime SLA for ingestion and alerting. Monitoring systems must survive the very datacenter failures they observe.
- Compression Efficiency: Achieve storage reduction versus raw 16-byte samples using Facebook's Gorilla Time-Series Compression Algorithm (Delta-of-deltas timestamp encoding + XOR float64 compression); the paper's production figure is , i.e. .
- Data Durability: (11 9s) durability across tiered storage backends.
2. Capacity & Scale Estimation (Back-of-the-Envelope Math)
Ingestion Scale & Network Bandwidth
- Sustained Ingest Rate: ().
- Peak Ingest Rate ( burst headroom): ().
- Raw Metric Point Size:
metric_namepointer (),timestamp(int64),value(float64), tag metadata pointers () .
- Peak Raw Network Ingress Bandwidth: Distributed across internal AWS Network Load Balancers (NLB) and an ECS Ingestion Gateway fleet.
Gorilla Compression Mathematics & Storage Derivations
Standard uncompressed time-series storage requires per point (). The Gorilla compression algorithm compresses this drastically:
- Timestamp Compression (Delta-of-Deltas):
- For regular scrape intervals (e.g., every 10 seconds), .
- Delta-of-deltas: .
- A delta-of-deltas of encodes as a single bit
0. - Over of all timestamps in production compress into a single
0bit.
- Float64 Value Compression (XOR Encoding):
- Successive floating-point measurements (e.g., CPU utilization fluctuating between and ) share identical IEEE 754 sign and exponent bits.
- XORing current and previous floats yields many leading and trailing zero bits.
- If value is unchanged (), encode as a single bit
0( of metrics in production). - If changed, encode control bits and only the meaningful XOR bits.
- Empirical Gorilla Compression Factor:
- Daily Compressed Ingestion Footprint:
Tiered Storage Footprint (Multi-Year Horizon)
- Hot Tier (In-Memory / NVMe on Amazon EKS - 2 Hours):
Series cardinality (the number people forget): at a 10-second scrape interval, means . Each series carries a Gorilla block header, a Roaring Bitmap posting entry per label, and a label-set string (roughly in total):
Sample data () plus index () with a safety margin needs across the hot tier. Because series are hash-sharded (not replicated) across nodes, a 3-node
r6i.4xlargecluster ( each, total) holds it with headroom; threer6i.2xlargenodes ( total) would not survive a single node loss. Durability during a node loss comes from the NVMe WAL replay plus the Kinesis retention window, not from RAM replication. - Warm Tier (Amazon Timestream / ClickHouse Columnar - 30 Days):
- Cold Tier (Downsampled Amazon S3 Parquet Lakehouse - 5 Years):
- Raw 1-second metrics are progressively downsampled:
- 1-minute rollups ( volume reduction) for days 31–180.
- 1-hour rollups ( volume reduction) for days 181–1,825 (5 years).
- Raw 1-second metrics are progressively downsampled:
3. High-Level Architecture & AWS Component Mapping
Synthesizing vector architecture diagram...
Why In-Memory Gorilla TSDB + Columnar Rollups? Slicing time-series telemetry into a hot in-memory tier (Gorilla compression holding raw 1-second samples for 2 hours) paired with an analytical warm tier (ClickHouse / Timestream holding 30-day 1-minute rollups) delivers sub- dashboard and alerting query response times while keeping storage footprint to .
Data Flow Walkthrough
- Telemetry Ingress: Agents collect metrics locally and push batched Protobuf payloads via persistent gRPC HTTP/2 connections to an internal AWS NLB.
- Consistent Partitioning: Ingest gateways inspect metric names and tags, compute a 64-bit MurmurHash3 shard key, and write records to Amazon Kinesis Data Streams.
- In-Memory Gorilla Ingestion: Stateful EKS worker pods consume from assigned Kinesis shards. Each series writes to an active 2-Hour Gorilla Block in RAM while appending raw bytes to an NVMe Write-Ahead Log (WAL) for crash recovery.
- Roaring Bitmap Indexing: Tags are indexed into an in-memory inverted index utilizing Roaring Bitmaps. Queries filtering on multiple labels execute bitwise
ANDintersections in microseconds. - Real-Time Alert Evaluation: The Alert Engine queries the Hot TSDB every 15 seconds. If an alert condition breaches its threshold across consecutive windows, the engine checks the Redis debounce ring (
SET NX EX 900). If new, a critical incident payload is published to Amazon SNS.
Concrete Step-by-Step Request Walkthrough: Ingesting & Querying a Metric Point
| Step # | Event / Action | Component State | Distributed Transition | Output / Response |
|---|---|---|---|---|
| 1 | OpenTelemetry Collector pushes batch of 1,000 metric points | Persistent gRPC stream over HTTP/2 connection | Ingestion gateway verifies tenant quota, validates schema & rate limits | Gateway acknowledges batch; returns HTTP 200 / gRPC OK |
| 2 | Gateway hashes metric identity | Key: cpu.util{env="prod", host="ip-10-0-1-4"} | MurmurHash3 maps time-series key to Kinesis Partition Shard 14 | Payload appended to Kinesis stream in latency |
| 3 | Gorilla TSDB consumes record | Shard 14 assigned to TSDB Worker Pod 03 | Worker appends point to local NVMe WAL ( wal_04.log) | NVMe sequential write verified; crash-recovery guarantee established |
| 4 | Gorilla Bit-Level Compression | Timestamp , Value | Delta-of-delta emits '0';Value unchanged from previous sample () emits '0' (1 bit) | Point stored in RAM slab consuming only physical memory |
| 5 | Roaring Bitmap inverted index update | Series ID assigned to unique label vector | Series ID inserted into __name__=cpu.utiland env=prod roaring bitmaps | Inverted index posting list ready for query in |
| 6 | Alert Rule Evaluation | Alert Engine evaluates:cpu.util{env="prod"} > 90% for 3m | Roaring Bitmap resolves candidate IDs; TSDB scans active 2-hour Gorilla blocks | Value ; evaluates clean with zero alerts fired |
4. API Interface Design & Wire Protocol
1. gRPC Ingestion Protocol (metrics_service.proto)
protobufsyntax = "proto3"; package hispeeddesign.metrics.v1; service MetricsIngestService { rpc PushMetrics (PushMetricsRequest) returns (PushMetricsResponse); } message PushMetricsRequest { string tenant_id = 1; repeated TimeSeries points = 2; } message TimeSeries { string metric_name = 1; map<string, string> labels = 2; repeated MetricSample samples = 3; } message MetricSample { int64 timestamp_ms = 1; double value = 2; } message PushMetricsResponse { uint32 points_ingested = 1; uint32 points_rejected = 2; string error_message = 3; }
2. Analytical Range Query REST API
httpPOST /v1/query_range HTTP/1.1 Host: tsdb-query.production.aws.internal Content-Type: application/json { "query": "avg(rate(http_requests_total{env=\"prod\", status=~\"5..\"}[5m])) by (service)", "start_epoch_sec": 1767222000, "end_epoch_sec": 1767225600, "step_seconds": 60 } Response: 200 OK Content-Type: application/json { "status": "success", "data": { "resultType": "matrix", "result": [ { "metric": { "service": "payment-api" }, "values": [ [ 1767222000, "0.012" ], [ 1767222060, "0.015" ], [ 1767222120, "0.084" ] ] } ] } }
5. Data Models & Storage Architecture
1. Gorilla Block Physical Memory Layout
Time-series points are organized into discrete 2-Hour In-Memory Slabs:
Each slab starts with a fixed header:
| Header field | Size |
|---|---|
| SeriesID | 4 B |
| BaseTimestamp | 8 B |
| BaseValue | 8 B |
| PointCount | 2 B |
After the header comes the compressed bitstream. Each point writes a timestamp code and then a value code; the first bits of each code (the prefix) say how many bits follow:
| Encoded bits | Meaning |
|---|---|
0 | Timestamp delta-of-delta == 0 |
0 | Value XOR == 0 (identical value) |
10 + 7 bits | Timestamp delta-of-delta in [-63, 64] |
110 + 9 bits | Timestamp delta-of-delta in [-255, 256] |
1110 + 12 bits | Timestamp delta-of-delta in [-2047, 2048] |
1111 + 32 bits | Timestamp delta-of-delta (anything else) |
1, 0, then the meaningful bits | Value XOR fits the previous leading/trailing-zero window |
1, 1, 5-bit leading-zero count, 6-bit length, then the bits | Value XOR with a new leading/length header |
2. Inverted Index Schema: Roaring Bitmaps
Labels are mapped to monotonically increasing 32-bit Series IDs:
| Label Key-Value Pair | Roaring Bitmap (Series ID Posting List) |
|---|---|
__name__="node.cpu.util" | {1, 4, 9, 12, 18, 42, 104, 1024, 8821} |
env="production" | {1, 4, 9, 15, 42, 104, 2048} |
region="us-east-1" | {1, 9, 18, 42, 104, 3012} |
| Bitwise AND Intersection | {1, 9, 42, 104} (Evaluated in CPU registers in ) |
3. DynamoDB Alert Rules Table (AlertRulesCatalog)
- Partition Key (
PK):TENANT#<tenant_id> - Sort Key (
SK):RULE#<rule_id> - Attributes:
query_expression,evaluation_interval_sec,for_duration_sec,severity,notification_sns_topic_arn,status.
Unlock Complete Architecture & Production Runbooks
You have explored the free architectural preview (~37%). Spend 1 Coin to unlock the remaining 6 production deep-dive sections for a full 24 hours.