Loading…
10 trillion samples a day: Scaling beyond traditional monitoring infra at Databricks
David Yuan, Yi Jin, Karan Bavishi, HC Zhu, Joey Beyda
- Source
- Databricks
- Published
- Added to Yomu
Summary
Databricks’ monitoring infrastructure now tracks 5 billion active timeseries in real time and ingests more than 10 trillion samples daily, exposing scalability, reliability, cost, and operability limits in its older stack. The replacement, Pantheon, is a fork of CNCF Thanos deployed across more than 160 instances in about 70 regions and three cloud providers; tiered storage, differentiated memory retention, isolated replicated Receive groups, multitenancy, and a custom control plane support automated scaling and recovery. Pantheon’s largest instance holds about 300 million in-memory timeseries and handles nearly 1,000 PromQL queries per second, while migration reduced annual cloud costs by millions and monitoring downtime by roughly five times. For high-cardinality troubleshooting, Hydra preserves raw metrics in Delta tables, exposes them through Grafana and SQL, and unifies metric semantics across aggregated and raw paths, with freshness improvements planned.
Context
Databricks’ previous monitoring stack was built for an order of magnitude lower scale and became a major reliability bottleneck. Growth in serverless and AI workloads sharply increased metric cardinality, while the infrastructure also had to operate reliably and consistently across roughly 70 cloud regions and three major cloud providers with minimal manual intervention.
Approach / What changed
Databricks built Pantheon, a Thanos fork with tiered storage, workload-specific memory retention, replicated and isolated Receive groups, multitenancy, and a custom control plane for rollouts, routing, autoscaling, and self-healing. It also built Hydra to preserve raw high-cardinality metrics in Delta tables, translate PromQL into SQL for Grafana, provide direct SQL access, and unify metric semantics across monitoring paths.
Takeaways
- Pantheon uses in-memory storage for recent data, on-disk storage for the last 24 hours, and object storage for older data, allowing compute to scale without rebalancing historical data across database nodes.
- Two Receive groups use different memory-retention policies for persistent and ephemeral workloads, while three isolated StatefulSets preserve quorum and allow parallel operational changes without affecting write availability.
- Hydra translates PromQL into SQL over Delta tables for Grafana and exposes those tables directly for SQL and notebook analysis, including joins with other enterprise datasets under shared security controls.