Taming 50 Billion Time Series: Operating Global-Scale Prometheus Dep... Orcun Berkem & Alan Protasio

Orcun Berkem, Alan Protasio

KubeCon + CloudNativeCon Europe 2025 · Session

Overview

This talk, presented by Orcun Berkem and Alan Protasio from AWS, delves into the intricate challenges and innovative solutions involved in operating a global-scale, Prometheus-compatible monitoring service on Kubernetes. Specifically, it details the architectural evolution and operational strategies behind AWS Managed Prometheus, a service designed to alleviate the operational burden of managing Prometheus infrastructure for users. The speakers, both deeply involved in open-source observability, highlight how they leverage and contribute to the Cortex project, an open-source, scalable, and highly available Prometheus-compatible system.

Watch on YouTube

Visual summary for Taming 50 Billion Time Series: Operating Global-Scale Prometheus Dep... Orcun Berkem & Alan Protasio by Orcun Berkem, Alan Protasio
Visual summary for Taming 50 Billion Time Series: Operating Global-Scale Prometheus Dep... Orcun Berkem & Alan Protasio by Orcun Berkem, Alan Protasio

Key moments

  1. 0:00 Welcome and talk agenda: global scale Prometheus challenges
  2. 1:45 AWS Managed Prometheus launch and open-source commitment
  3. 3:00 Understanding the inherent limitations of standalone Prometheus
  4. 4:20 Introduction to Cortex and its microservices architecture
  5. 6:05 Cortex's horizontal scalability and data rebalancing
  6. 8:05 Cortex's robust multi-tenancy and isolation capabilities
  7. 9:05 Achieving long-term data retention with Cortex

Taming 50 Billion Time Series: Operating Global-Scale Prometheus Deployments on Kubernetes

Speakers: Orcun Berkem, Principal Engineer, AWS; Alan Protasio, Software Engineer, AWS, Cortex Maintainer

Conference: KubeCon EU

YouTube: https://www.youtube.com/watch?v=OqLpKJwKZlk

Overview

This talk, presented by Orcun Berkem and Alan Protasio from AWS, delves into the intricate challenges and innovative solutions involved in operating a global-scale, Prometheus-compatible monitoring service on Kubernetes. Specifically, it details the architectural evolution and operational strategies behind AWS Managed Prometheus, a service designed to alleviate the operational burden of managing Prometheus infrastructure for users. The speakers, both deeply involved in open-source observability, highlight how they leverage and contribute to the Cortex project, an open-source, scalable, and highly available Prometheus-compatible system.

The core of the presentation addresses the inherent limitations of a standalone Prometheus deployment—its monolithic nature, single-machine constraint, and lack of built-in redundancy and multi-tenancy support—when confronted with the demands of an enterprise-grade, global service. By adopting and enhancing Cortex, AWS has built a robust platform capable of handling tens of billions of time series, ensuring high availability, multi-tenancy, and significant blast radius reduction. The talk is particularly relevant for engineers, architects, and site reliability professionals grappling with the complexities of large-scale observability, distributed systems, and Kubernetes operations in multi-tenant environments.

The speakers meticulously outline the journey from identifying Prometheus's scaling limitations to implementing a sophisticated cellular architecture that underpins AWS Managed Prometheus. They cover critical aspects such as deployment safety, capacity management, tenant isolation, and advanced testing methodologies. This deep dive offers invaluable insights into real-world solutions for managing massive data volumes, ensuring service resilience, and maintaining a rapid pace of development and deployment within a highly distributed, cloud-native ecosystem.

Background

▶ Watch: Welcome and talk agenda: global scale Prometheus challenges (0:00)

The journey into global-scale observability, as presented by Berkem and Protasio, is rooted in the evolution of the Prometheus ecosystem. Prometheus, a purpose-built open-source time series database with a powerful query language, PromQL, was started in 2012 and open-sourced in 2015. It quickly gained traction for its robust monitoring capabilities, particularly within the Kubernetes ecosystem, thanks to its wide range of exporters and strong community adoption. However, for a managed service like AWS Managed Prometheus, launched in 2020, Prometheus's inherent design presented significant challenges.

The primary limitation of Prometheus is its monolithic architecture, operating as a single process on a single machine. This design leads to several critical issues at scale:

  • Resource Contention: Ingest-heavy workloads compete with query execution and alerting rules for the same host resources, leading to performance degradation.
  • No Redundancy: A single point of failure means no built-in high availability.
  • Retention Limits: Data retention is constrained by the local disk capacity of a single host, making long-term storage (multiple years) impractical.
  • Lack of Multi-tenancy: Prometheus does not natively isolate data or resources between different users, making it unsuitable for a shared, managed service without significant custom engineering.

Recognizing these limitations, AWS explored alternative solutions and ultimately selected Cortex, a CNCF project that emerged shortly after Prometheus. Cortex was designed specifically to address Prometheus's scaling challenges by breaking down its monolithic structure into a distributed system of microservices. Cortex ingeniously imports the Prometheus codebase as a library, ensuring full backward compatibility with Prometheus and enabling easy backporting of new features like native histograms.

Cortex's architecture provides several foundational advantages:

  • Horizontal Scalability: Different microservices can be scaled independently based on workload (e.g., read-heavy vs. write-heavy), freeing the system from the constraints of a single host.
  • High Availability: Data replication across nodes and availability zones, combined with quorum operations, eliminates single points of failure.
  • Multi-tenancy: Cortex offers built-in isolation for tenant data, along with limits and quotas and shuffle sharding to mitigate "noisy neighbor" issues.
  • Long-Term Storage: By offloading historical data to object storage like Amazon S3 or Google Cloud Storage (GCS), Cortex overcomes local disk limitations, supporting data retention for years.
  • Auto-Balancing: Cortex automatically rebalances load and re-shards data across nodes during scaling activities using consistent hashing, minimizing disruption.

These capabilities made Cortex an ideal foundation for AWS Managed Prometheus, allowing AWS to offer a serverless, scalable, highly available, and secure Prometheus-compatible monitoring service that integrates seamlessly with other AWS services and supports long-term data retention.

Key Findings

▶ Watch: Understanding the inherent limitations of standalone Prometheus (3:00)

The central challenge identified in operating a global-scale monitoring service, even with a distributed system like Cortex, is that a single, infinitely growing cluster inevitably encounters hidden scaling cliffs and non-linear scaling behaviors. These issues, stemming from communication overhead and coordination complexities in very large distributed systems, can lead to a state where adding more hosts actually degrades performance. This realization led to the development and adoption of a cellular architecture.

The cellular architecture is the talk's paramount finding, providing a robust solution for global-scale operations. Instead of one massive Cortex cluster, AWS deploys multiple, self-contained Cortex clusters, referred to as "cells," within each region. Each cell is an independent deployment, including its own S3 bucket, authorization infrastructure, and Kubernetes cluster. Tenants are assigned to specific cells, providing several critical benefits:

  • Enhanced Scalability: New cells can be provisioned to handle increased load or new clients, allowing for indefinite horizontal scaling across the service.
  • Improved Deployment Safety: Changes are rolled out cell-by-cell, localizing potential issues to a small subset of customers, rather than impacting an entire region or the global service.
  • Blast Radius Reduction: If a tenant or an internal issue causes a cell to become unhealthy, the impact is confined to that specific cell and its assigned tenants, preventing widespread outages.
  • Decreased Mean Time To Recover (MTTR): Operators can focus on a smaller, isolated fleet and customer subset during incidents, significantly accelerating diagnosis and recovery.

Beyond the cellular architecture, the speakers highlighted several other key findings crucial for operating at this scale:

  • Zone-Aware Controllers for StatefulSets: During Kubernetes deployments, standard Pod Disruption Budgets don't consider availability zones. This can lead to multiple ingesters (which are StatefulSets) being taken down across different AZs simultaneously, causing service disruption. AWS developed and open-sourced zone-aware controllers that prioritize replacing nodes within a single AZ before moving to the next, drastically speeding up deployments and minimizing customer impact.
  • Proactive Host Health Monitoring: Traditional health checks often fail to detect "gray failures" – subtle issues like network degradation or disk errors that don't immediately crash a host but impair its performance. AWS implements host health monitoring using multiple signals (e.g., 5xx error counts, latency) to identify and proactively replace unhealthy nodes, even in cases of false positives, to ensure continuous service availability.
  • Adoption of Carpenter: This Kubernetes cluster auto-scaling project offers significant advantages for global-scale operations. Carpenter provides better transparency and debuggability during deployments, faster scaling characteristics, and the ability to pick multiple instance types, mitigating capacity constraints on specific host types.
  • Future Architectural Enhancements: The team is actively working on partitions within Cortex to further improve resiliency against multiple ingester failures across different AZs, by assigning series to partitions which are then assigned to ingesters. They are also exploring the Parquet format for long-term storage to reduce operational toil associated with store gateways and potentially enhance query performance.

These findings collectively demonstrate a sophisticated approach to building and operating a highly resilient, scalable, and manageable observability platform in a demanding cloud environment.

Technical Deep Dive

▶ Watch: Introduction to Cortex and its microservices architecture (4:20)

The technical foundation of AWS Managed Prometheus is Cortex, an open-source system designed to scale Prometheus horizontally. Cortex breaks down the monolithic Prometheus server into distinct microservices, each with a specialized role, which can be scaled independently. This architecture is broadly divided into three paths:

  1. Write Path: Involves Distributors and Ingesters. Distributors receive incoming Prometheus metrics, fan them out to multiple ingesters for high availability (replicating data across nodes and AZs), and apply initial processing. Ingesters temporarily store "hot" data in memory and on local disk, then periodically flush blocks of data to long-term object storage.
  2. Read Path: Comprises Store Gateways, Queries, and Query Frontends. Query Frontends handle incoming PromQL queries, potentially fanning them out to multiple Query services. Queries retrieve data from ingesters (for recent data) and store gateways (for historical data). Store Gateways are responsible for reading data blocks from the long-term object storage (e.g., S3) and serving them to queries.
  3. Alerts and Ruler Path: Includes the Alert Manager and Rulers. Rulers evaluate configured alerting and recording rules against incoming data, while Alert Manager handles the routing and deduplication of alerts.

A crucial aspect of Cortex's design for a managed service is its robust multi-tenancy support. Cortex isolates data from different tenants, ensuring that one tenant's data is not accessible by another. Beyond data isolation, it implements noise neighbor mitigations through:

  • Limits and Quotas: Each tenant is allocated specific limits on metrics ingested, query complexity, and resource usage, preventing a single tenant from monopolizing shared resources.
  • Shuffle Sharding: This strategy minimizes the overlap of resources shared between tenants. Even if a "bad tenant" overloads a few hosts, other tenants sharing those hosts are less likely to be impacted due to redundancy and the limited sharing surface.

For long-term storage, Cortex keeps only "hot" data on local disk for a few hours. This data is then shipped to highly durable and scalable object storage like S3. This mechanism allows for data retention periods extending from days to multiple years, a critical requirement for many customers.

The fundamental innovation to scale Cortex itself globally is the cellular architecture. This architecture is composed of:

  • Cell Router: A minimal, thin layer acting as the single point of entry for customers. Its sole responsibility is to route incoming requests to the correct cell based on tenant assignment. This layer is transparent to the customer.
  • Cell: Each cell is a complete, self-contained deployment of Cortex, along with all its dependencies: its own Kubernetes cluster, dedicated S3 bucket, authorization systems, and other necessary infrastructure. This ensures maximum isolation.
  • Control Plane: This layer is responsible for administrative tasks, including assigning tenants to cells, provisioning new cells when needed, and deprovisioning cells that are no longer required.

Cell Capacity Management is vital. Each cell operates within a predefined "safe limit," determined through extensive load testing. Two thresholds govern a cell's lifecycle:

  • Scaling Threshold: When a cell's capacity approaches this threshold, the control plane stops assigning new tenants to it. Instead, new tenants are routed to newly provisioned cells.
  • Hard Limit: An absolute maximum capacity that a cell should never exceed.

To prevent existing tenants from growing beyond the hard limit, an automated controller within each cell manages limits and quotas. This controller continuously monitors customer usage, auto-grants limit increases when deemed safe (i.e., within the cell's hard limit), and reclaims unused limits to free up capacity. This controller also pre-provisions nodes as a cell grows organically. When a new cell is created, it starts with a minimal footprint and scales up as needed. Before granting a customer a limit increase that requires more resources, the controller ensures that the necessary nodes are pre-provisioned and ready.

A key operational challenge addressed is cell migration. If an existing tenant within a cell requires more capacity than the cell can safely provide, or if a cell becomes unbalanced, a migration is triggered. This involves:

  1. Creating a new cell or identifying an existing cell with sufficient free capacity as the target.
  2. Mirroring traffic (both write and read) between the source and target cells for a period to ensure data consistency.
  3. Migrating configurations and historical data. Crucially, historical data, being in the source cell's S3 bucket, must be transferred to the target cell's S3 bucket, as cells do not share object storage.
  4. Once data is consistent, stopping traffic mirroring and routing all operations to the new cell. This process is designed to be seamless for the customer, with the only potential impact being occasional duplicate alerts during the final transition.

For releasing changes across this complex global infrastructure, AWS employs deployment waves. The process ensures safety while respecting cellular and regional boundaries:

  1. Pre-production: Unit and integration tests in beta environments.
  2. Production Waves:
  • Gamma Environment: Changes are first deployed to two cells in a small region to test Cortex and the cell router. Canaries run 24/7, and signals are monitored for several hours.
  • First Production Wave (Small Region): A single cell in a small region receives the change. Automated alarms trigger rollbacks if issues arise, localizing impact.
  • Remaining Cells (Small Region): Once confidence is gained, the change rolls out to other cells in that region.
  • Next Big Wave (Multiple Regions): The scope expands to two or more regions, again starting with gamma deployments, then single cells per region, and finally all remaining cells. This exponential expansion continues until all regions are updated. The entire process for a single change can take "on the order of days."

Rigorous testing underpins this process, including fuzziness and correctness tests for PromQL queries, comparing results against a vanilla Prometheus instance to prevent regressions. Fault injection is also used between production waves to validate the system's resilience mechanisms against various impairments.

Finally, managing deployments within a Kubernetes cluster also requires specialized solutions. For StatefulSets like Cortex ingesters, standard Kubernetes Pod Disruption Budgets don't consider availability zones. Taking down two ingester pods from different AZs simultaneously can cause service disruption due to how series are sharded. To counter this, AWS developed and open-sourced zone-aware controllers. These controllers deploy changes to ingesters in a zone-by-zone fashion, starting with a small number of pods (1, then 2, then 4), exponentially growing within one AZ before moving to the next. This significantly speeds up deployments while preventing cross-AZ disruption.

Demo / Proof of Concept

▶ Watch: Cortex's robust multi-tenancy and isolation capabilities (8:05)

The talk primarily focused on architectural design and operational strategies rather than a live demonstration or a specific proof of concept. The speakers detailed the implementation and benefits of their global-scale Prometheus deployment using Cortex and a cellular architecture, illustrating the "how" through diagrams and explanations of their deployment processes and system behaviors.

Defensive Implications

▶ Watch: Achieving long-term data retention with Cortex (9:05)

Operating a global-scale monitoring service like AWS Managed Prometheus offers critical lessons for defenders and system architects:

  • Embrace Distributed Architectures for Scalability and Resilience: The move from monolithic Prometheus to a microservices-based system like Cortex is a fundamental defensive strategy. It allows for independent scaling of components, preventing single points of failure and resource contention, which inherently improves the system's ability to withstand various loads and failures.
  • Implement Strong Multi-Tenancy Controls: For any shared service, robust tenant isolation is paramount. Cortex's approach with data isolation, limits and quotas, and shuffle sharding demonstrates how to prevent "noisy neighbor" issues and ensure that one compromised or misbehaving tenant does not impact others. Defenders should ensure their multi-tenant systems have similar controls.
  • Adopt Cellular or Multi-Cluster Strategies for Blast Radius Reduction: The cellular architecture is a powerful pattern for limiting the impact of failures. By containing incidents within a single cell, an organization can prevent global outages, significantly reduce Mean Time To Recover (MTTR), and improve overall service availability. This strategy is applicable beyond monitoring to any critical, large-scale distributed system.
  • Prioritize Deployment Safety through Phased Rollouts and Automated Rollbacks: The deployment waves strategy, with its emphasis on canary testing, gamma environments, and automated rollbacks, is a best practice for safe, large-scale software delivery. Defenders should advocate for similar cautious, incremental deployment processes to minimize the risk of introducing vulnerabilities or regressions into production.
  • Develop Zone-Aware Deployment Mechanisms for Stateful Workloads: For StatefulSets running across multiple availability zones, standard Kubernetes deployment strategies can inadvertently cause service disruption. Custom solutions like zone-aware controllers are essential to ensure that critical stateful components maintain high availability during upgrades and failures. This highlights the need for deep understanding and customization of cloud-native infrastructure.
  • Proactively Detect and Remediate "Gray Failures": Relying solely on basic health checks is insufficient for large-scale operations. Implementing host health monitoring that uses multiple signals (e.g., 5xx counts, latency) to detect subtle performance degradations—gray failures—and trigger proactive host replacement is crucial. This pre-emptive approach prevents minor issues from escalating into major outages.
  • Leverage Advanced Cloud-Native Tools for Infrastructure Management: Tools like Carpenter for Kubernetes cluster autoscaling can improve operational efficiency, provide better transparency during scaling events, and offer greater flexibility in instance type selection. Defenders should explore how such tools can enhance the resilience and manageability of their infrastructure.
  • Plan for Long-Term Data Retention with Scalable Storage: Utilizing object storage (like S3) for historical data is a cost-effective and scalable solution for managing massive volumes of time series data over extended periods. This ensures that critical forensic and analytical data is available for compliance, security investigations, and long-term trend analysis.
  • Continuously Innovate for Resilience: The ongoing work on partitions for improved fault tolerance and the exploration of the Parquet format for storage optimization demonstrate a commitment to continuous improvement. Defenders should foster a culture of architectural evolution to address emerging challenges and enhance system resilience over time.

Key Takeaways

  • Prometheus's monolithic nature necessitates distributed solutions for global scale: While powerful, vanilla Prometheus lacks the inherent scalability, high availability, and multi-tenancy features required for a large-scale managed service, leading AWS to adopt and enhance the Cortex project.
  • Cellular architecture is crucial for managing complexity and limiting blast radius: Deploying multiple, self-contained Cortex "cells" per region, each with its own infrastructure, significantly improves scalability, deployment safety, and dramatically reduces the impact of localized failures or misbehaving tenants.
  • Sophisticated deployment strategies are essential for safe, large-scale releases: AWS employs deployment waves with canary testing, fault injection, and automated rollbacks across gamma and production environments, along with specialized zone-aware controllers for StatefulSets, to ensure changes are rolled out safely and efficiently across tens of thousands of hosts.
  • Proactive monitoring and automated remediation are vital for detecting "gray failures": Beyond basic health checks, host health monitoring uses multiple signals like 5xx counts and latency to identify and automatically replace subtly failing nodes, preventing service disruption before it becomes critical.
  • Continuous architectural evolution drives resilience and operational efficiency: Ongoing work on partitions to enhance fault tolerance against multiple ingester failures and exploring the Parquet format for long-term storage aims to further reduce operational toil and improve system performance in the face of ever-increasing scale.

About the Speaker(s)

Orcun Berkem is a Principal Engineer at AWS, where his work focuses on open-source observability solutions. He plays a key role in developing and operating the global-scale infrastructure discussed in the talk, addressing complex challenges related to availability, multi-tenancy, and blast radius reduction.

Alan Protasio is a Software Engineer at AWS and a dedicated Cortex maintainer. His contributions to the Cortex project are integral to its evolution and its adoption within AWS Managed Prometheus. He brings deep expertise in distributed systems and the Prometheus ecosystem to the team.

Reviews

Dr. Zero (Offensive Security Researcher) — MUST SEE

This talk delivers a masterclass in operating a global-scale Prometheus-compatible monitoring service. The speakers, deeply involved in the Cortex project and AWS Managed Prometheus, meticulously detail the architectural evolution from Prometheus's inherent limitations to a sophisticated cellular architecture. They present novel solutions for multi-tenancy, blast radius reduction, deployment safety, and proactive failure detection, all while leveraging and contributing to open-source projects. This is a rare look into the real-world engineering challenges and solutions behind a critical cloud-native service.

Heather Calloway (CISO) — STRONG ACCEPT

This KubeCon talk meticulously details the architectural and operational evolution of AWS Managed Prometheus, moving from a monolithic system to a highly resilient, cellular architecture built on Cortex. While deeply technical, the presentation offers critical insights into managing global-scale distributed systems, providing a robust blueprint for blast radius reduction, multi-tenancy, and deployment safety. Its focus on institutional realism and engineered resilience makes it highly relevant for security leaders grappling with operational risk and accountability in cloud-native environments.

→ Top-rated talks at KubeCon + CloudNativeCon Europe 2025

All talks from KubeCon + CloudNativeCon Europe 2025