When we started CloudMonitor in 2022, our metrics pipeline handled about 50 million data points per day. Today, that number exceeds 20 trillion. Scaling a real-time observability system by six orders of magnitude forced us to rethink nearly every architectural decision we made in the early days.
The Original Architecture
Version 1 was straightforward: agents pushed metrics to a central HTTP API, which wrote to a time-series database (TimescaleDB), and a React frontend queried it directly. This worked beautifully for our first 50 customers. Then it didn't.
Ingestion: From Push to Buffered Pipeline
The first bottleneck was ingestion. At peak traffic, the HTTP API would hit connection limits, causing agents to buffer and retry, which compounded the problem. We replaced the synchronous push model with an asynchronous pipeline: agents write to local Kafka topics, and a fleet of ingestion workers consume, validate, and route metrics to storage.
This single change improved ingestion throughput by 40x and gave us backpressure handling for free — when storage slows down, Kafka buffers rather than dropping data.
Storage: Specialized Clusters
TimescaleDB served us well, but at 20 trillion metrics per day, the write amplification from PostgreSQL's WAL became unsustainable. We migrated to a tiered storage model: hot data (last 24 hours) lives in an in-memory columnar store, warm data (24 hours to 30 days) in a custom LSM-based storage engine optimized for time-series workloads, and cold data in S3-compatible object storage with Parquet encoding.
Query Engine: Pre-Aggregation and Materialized Views
Users expect dashboard queries to return in under 500ms, regardless of the time range. Pre-computing rollups at 1-minute, 5-minute, 1-hour, and 1-day granularities ensures that query cost is proportional to the number of data points displayed, not the total volume stored.