Pushing the Limits of Prometheus at Etsy - Chris Leavoy, Etsy & Bryan Boreham, Grafana Labs
Chris Leavoy, Etsy, Bryan Boreham, Grafana Labs
KubeCon + CloudNativeCon Europe 2025 · Session
Overview
This article delves into the intricate challenges and innovative solutions Etsy encountered while operating one of the largest single Prometheus servers in the industry. Presented by Chris Leavoy from Etsy and Bryan Boreham from Grafana Labs, the talk provides a deep dive into the complexities of scaling Prometheus vertically to its absolute limits, revealing unexpected bottlenecks and offering practical strategies for optimization. The speakers share invaluable lessons learned from dealing with millions of metrics, high churn environments, and the critical need for robust observability during peak periods like Black Friday.
Key moments
- 0:00 Introduction, speakers, and Prometheus architecture.
- 2:00 Etsy's 30 Prometheus stacks, 600M series.
- 3:00 Why Etsy chose vertical scaling and 'boring tech'.
- 4:00 Server overloaded, dropping data, and metric surges.
- 5:00 Deployment churn causes massive metric spikes.
- 6:00 Hardware upgrades fail; bottleneck discovered, production benchmarking.
- 8:00 Grafana Labs' Bryan Boreham gets involved.
- 8:00 Diagnosing severe Prometheus remote write lag.
Pushing the Limits of Prometheus at Etsy
Speakers: Chris Leavoy, Observability Engineer, Etsy; Bryan Boreham, Distinguished Engineer, Grafana Labs
Conference: KubeCon EU
YouTube: https://www.youtube.com/watch?v=IWrd-pSojqg
Overview
This article delves into the intricate challenges and innovative solutions Etsy encountered while operating one of the largest single Prometheus servers in the industry. Presented by Chris Leavoy from Etsy and Bryan Boreham from Grafana Labs, the talk provides a deep dive into the complexities of scaling Prometheus vertically to its absolute limits, revealing unexpected bottlenecks and offering practical strategies for optimization. The speakers share invaluable lessons learned from dealing with millions of metrics, high churn environments, and the critical need for robust observability during peak periods like Black Friday.
The discussion highlights Etsy's pragmatic approach of "choosing boring technology" and exhausting the potential of a single-instance architecture before resorting to more complex distributed solutions. This philosophy led them to push Prometheus far beyond typical deployments, uncovering fundamental performance limitations within its core design, particularly concerning remote write, write-ahead log operations, and Go runtime garbage collection. The collaboration between Etsy's operational expertise and Grafana Labs' deep Prometheus knowledge ultimately led to significant optimizations, many of which have been upstreamed to benefit the broader Prometheus community.
The talk is crucial for anyone operating Prometheus at scale, or those considering its deployment in high-volume, dynamic environments. It demystifies performance issues often attributed to high cardinality and instead points to architectural and runtime bottlenecks that can severely impact even well-resourced instances. By detailing their journey from frequent server crashes and data loss to a stable, high-performing system, Leavoy and Boreham offer a blueprint for diagnosing and resolving similar challenges, emphasizing the importance of meticulous tuning and understanding the underlying mechanisms of the monitoring stack.
Background
▶ Watch: Introduction, speakers, and Prometheus architecture. (0:00)
Etsy, a global e-commerce marketplace, operates a sophisticated observability infrastructure that includes 30 isolated Prometheus stacks, each dedicated to monitoring specific system boundaries. Collectively, these stacks handle a peak of around 600 million time series, processing approximately 5 million samples per second during busy periods. Each stack also writes a copy of its data to a central location in Grafana Cloud for long-term retention and centralized querying. The focus of this talk is on Etsy's largest Prometheus stack, which monitors millions of metrics related to searches on etsy.com and has been operational for about eight years, accumulating a substantial number of alerts and recording rules.
Etsy's engineering culture has historically favored vertical scaling and simplicity, adhering to the principle of "choosing boring technology." This meant exhausting all possibilities with a single-instance architecture before introducing the complexity of sharding or federation. This strategy, while initially effective, led them to a critical juncture where their largest Prometheus server could no longer scale vertically. The server was frequently overloaded, experiencing dropped data, skipped recording rules, and missed scrapes, particularly during periods of high churn caused by Etsy's Kubernetes-based auto-scaling and blue-green deployments. Each deployment could involve 50 million metrics stopping and another 50 million coming online, with an overlapping period, creating significant stress on the system.
The impending Black Friday and Cyber Monday season, Etsy's busiest time of year, added urgency to the situation. Despite upgrading to the largest available Google Cloud instance (an ultr instance with 4 terabytes of RAM, up from 2 terabytes), the issues persisted, indicating a bottleneck beyond mere hardware capacity. This led to a crucial decision: deploying a third replica of the largest stack to conduct side-by-side benchmarks in production, enabling rapid iteration and troubleshooting without jeopardizing the live service. It was during this phase that they discovered the server performed better with significantly less RAM (1 terabyte) than the initial 4 terabytes, signaling deeper, non-memory-related performance issues that required a detailed investigation.
Key Findings
▶ Watch: Why Etsy chose vertical scaling and 'boring tech'. (3:00)
Etsy's journey to stabilize their largest Prometheus instance uncovered several critical bottlenecks and led to key findings that challenged conventional scaling wisdom:
- Remote Write Lag: The most immediate and alarming issue was Prometheus remote write lagging by 5 to 10 minutes, persisting for hours. This meant critical monitoring data was not reaching Grafana Cloud in a timely manner, severely impacting centralized observability. The
desired_shardsmetric was found to be highly misleading, suggesting an impossibly high number of shards that exacerbated performance problems rather than solving them. - Network Throughput Limitations: An unexpected discovery was a 360 megabits per second (Mbps) egress limit for single flows from Google Compute Engine instances to destinations outside the Virtual Private Cloud (VPC). Prometheus, by default, was sending remote write data over a single HTTP/2 socket, hitting this hard limit and causing timeouts, even when more bandwidth was theoretically available.
- Write-Ahead Log (WAL) Bottleneck: Despite significant parallelism in scraping and remote write components, the core write-ahead log (WAL) was identified as a single-threaded bottleneck. All incoming data had to pass through this sequential operation before being written to disk, and subsequently read for remote write, effectively becoming the choke point for the entire system.
- Compaction Overheads in High-Churn Environments: Prometheus's default compaction mechanisms, designed to optimize historical data storage, proved detrimental in Etsy's high-churn environment. Frequent blue-green deployments caused head compaction to fail due to internal index size limits (specifically, exceeding 64 gigabytes for a 2-hour block), leading to an infinitely growing WAL and requiring manual intervention and data loss. Additionally, restarts could take over an hour, exceeding Kubernetes timeouts.
- Go Runtime Garbage Collection Interference: Even after addressing other issues, the system exhibited periodic throughput dips every two minutes. This was directly linked to the Go runtime's garbage collection (GC) cycles, which would briefly starve the remote write process, causing performance oscillations. The default GC behavior, especially when compounded by large memory allocations (like during compaction), significantly impacted overall stability.
- Diminishing Returns from Hardware: The initial assumption that "more hardware is better" was debunked. The server performed worse with 4 terabytes of RAM than with 1 terabyte, and reducing CPU cores from 128 to 50 actually improved performance. This highlighted that resource saturation was not the primary problem; rather, it was inefficient resource utilization and architectural bottlenecks.
Technical Deep Dive
▶ Watch: Deployment churn causes massive metric spikes. (5:00)
The investigation into Prometheus's performance at Etsy revealed several deep technical insights, primarily focused on optimizing the Go runtime and Prometheus's internal mechanisms. The standard Prometheus architecture involves exporters sending data, which is then scraped by the Prometheus server and stored in its time series database (TSDB) – first in memory, then on disk. Data is queried using PromQL and typically visualized in tools like Grafana. The critical aspect for this talk is that the entire TSDB, scraping, and remote write logic runs as a single process, making vertical scaling a natural first approach.
Remote Write Tuning and Throughput Calculation
Etsy's initial problem manifested as severe remote write lag, with data taking 5-10 minutes to reach Grafana Cloud. Prometheus attempts to balance data ingestion and remote write by dynamically adjusting shards, but the prometheus_remote_storage_desired_shards metric proved misleading, suggesting thousands of shards. This was a "fantasy number" as it assumed infinitely fast sending and receiving.
The team eventually hardcoded min_shards and max_shards to 100, a value chosen relative to the server's CPU count. This was significantly less than the thousands suggested by the metric. They also increased batch_size and buffer_size to match. A crucial step was calculating the maximum theoretical throughput under worst-case latency:
- Data rate: 30 million series scraped every 15 seconds, equating to 2 million samples per second.
- Latency: Assuming a worst-case roundtrip latency of 500 milliseconds (0.5 seconds).
- Batches per second per shard: With 500ms latency, each shard can send two batches per second.
- Samples per batch: Each batch was configured for 20,000 samples.
- Theoretical throughput: 100 shards 2 batches/second/shard 20,000 samples/batch = 4 million samples per second.
This calculation showed ample headroom at 500ms latency. However, doubling the latency to 1 second would halve the throughput to 2 million samples per second, indicating a critical threshold. This methodical approach allowed them to configure remote write for optimal performance, moving away from dynamic, often over-aggressive, shard adjustments.
Network Egress Bottleneck
During the remote write investigation, network bandwidth charts occasionally showed timeouts, and egress peaked at 360 Mbps, even when more throughput was expected. This led to the discovery of a Google Cloud limitation: a single flow egressing from a compute instance to a destination outside its VPC cannot exceed approximately 360 Mbps. A netstat check revealed Prometheus was indeed using a single socket for remote write. The solution was to disable HTTP/2 in the Prometheus configuration. This forced Prometheus to use multiple HTTP/1.1 connections, bypassing the single-flow limit and allowing full utilization of the underlying network capacity.
Write-Ahead Log (WAL) as the Central Bottleneck
Even with remote write and network issues addressed, persistent lag remained. Profiling, while useful for average program performance, didn't pinpoint the issue. Instead, Go execution tracing provided the "smoking gun." The trace showed a solid yellow line at the top, representing the single goroutine reading the WAL, while other goroutines (remote write shards) spent most of their time in whitespace, waiting for data.
This confirmed that the write-ahead log (WAL), a fundamental design pattern for data resiliency in databases like Prometheus, was the central bottleneck. While excellent for recovery, in a massive single-instance setup, the WAL's sequential nature meant all incoming scrapes and outgoing remote writes were bottlenecked on this single disk operation. The team's immediate goal shifted from optimizing parallel operations to making this single-threaded WAL operation as fast as possible.
Compaction Strategy for High-Churn Environments
Prometheus's TSDB compaction process, which aggregates smaller in-memory data blocks into larger, more efficient on-disk blocks (e.g., from 2-hour blocks to 6-hour, then 18-hour blocks), became problematic for Etsy.
- Restart times: Default 2-hour blocks meant Prometheus had to rebuild two hours of data from the WAL on restart. During peak times, this could exceed Kubernetes' 1-hour startup timeout, leading to perpetual restarts unless the WAL was manually deleted, resulting in data loss.
- Head compaction failures: Etsy's frequent blue-green deployments and auto-scaling generated immense metric churn. This churn caused the head compaction (from memory to disk) to fail because the resulting 2-hour blocks exceeded an internal index size limit of 64 GB. This would leave Prometheus with an infinitely growing WAL that could not be compacted, again necessitating manual intervention and data loss.
To resolve this, Etsy drastically altered their compaction strategy:
- They set both
min_block_durationandmax_block_durationto 1 hour. - This configuration forces Prometheus to write a continuous stream of 1-hour blocks, effectively disabling the subsequent, larger historical compactions. The
min_block_durationsetting also ensures series clear out of the head block faster.
This strategy has trade-offs: many smaller blocks mean historical queries are more expensive. However, Etsy mitigates this by retaining only 9 days of data (216 1-hour blocks) locally, compared to the default 15 days. Crucially, expensive historical queries are offloaded to Grafana Cloud, and aggressive query limits (e.g., query_max_samples at 50 million) are set on the local Prometheus instance to prevent overload.
Go Runtime Garbage Collection Tuning
The final piece of the puzzle involved the Go runtime garbage collection (GC). The throughput graphs showed regular dips every two minutes, directly corresponding to GC cycles. The Go GC operates on a "sawtooth" model: the heap grows, and when it reaches current_heap_size + 75% (controlled by GOGC setting, default 100%), a GC cycle runs, reducing memory usage.
Compaction, especially the larger historical compactions, would temporarily increase the baseline memory usage, leading to even larger and more impactful GC cycles. The solution involved a combination of right-sizing the server and explicit Go runtime tuning:
- Right-sizing: The server was downsized from 128 cores to 50 cores and from 4 terabytes to 1 terabyte of RAM. This counter-intuitive step was critical because more memory allowed the heap to grow larger before GC, leading to longer, more disruptive GC pauses.
GOMEMLIMIT: SettingGOMEMLIMIT=1TBcaps the Go process's memory usage. When the heap approaches this limit, the Go runtime performs GC more aggressively. This, combined with disabled compaction and limited queries, meant memory usage was predictable.GOGC=off: By settingGOGC=off, the default 2-minute periodic GC cadence was disabled. GC was now primarily triggered by theGOMEMLIMITand the natural memory allocation patterns, leading to GC cycles occurring every 10 minutes instead of every 2 minutes.
This combination significantly reduced the frequency and impact of GC pauses, allowing the remote write process to catch up and maintain consistent throughput. Many of these optimizations, particularly around Go runtime and TSDB efficiency, were contributed back to the open-source Prometheus project, benefiting the wider community.
Demo / Proof of Concept
▶ Watch: Hardware upgrades fail; bottleneck discovered, production benchmarking. (6:00)
While no explicit "demo" in the traditional sense was presented, Etsy's entire troubleshooting process served as a live, production-scale proof of concept. The deployment of a third, identical Prometheus replica in production allowed for rapid, side-by-side benchmarking and iteration on configuration changes. This real-world experimentation, directly on production traffic, provided irrefutable evidence of the effectiveness of each optimization, from remote write tuning to Go runtime adjustments, under the most demanding conditions leading up to Black Friday.
Defensive Implications
▶ Watch: Diagnosing severe Prometheus remote write lag. (8:00)
The detailed account of Etsy's Prometheus journey offers a robust checklist of defensive strategies for anyone operating or planning to deploy Prometheus at scale:
- Automate Metric Explosion Detection: Implement automated checks to catch code changes that could lead to sudden, massive increases in metric cardinality or volume before they are deployed to production. This proactive measure can prevent server overloads and outages.
- Set Scrape Limits: The default Prometheus configuration has no upper bound on the number of metrics it will accept, which is a significant "foot gun." Defenders should configure explicit
scrape_limitsto prevent a single misbehaving exporter or application from overwhelming the Prometheus server. - Optimize Compaction for Churn: In environments with high metric churn (e.g., frequent auto-scaling, blue-green deployments), consider using smaller
min_block_durationandmax_block_durationsettings (e.g., 1 hour) to disable larger historical compactions. This can prevent head compaction failures, reduce restart times, and avoid an infinitely growing write-ahead log. Understand the trade-offs regarding historical query performance. - Tune Remote Write Carefully: Ignore the often-misleading
prometheus_remote_storage_desired_shardsmetric. Instead, calculate and configuremin_shardsandmax_shardsbased on CPU cores, destination latency, and required throughput. Remember to adjustbatch_sizeandbuffer_sizeaccordingly. - Right-Size Servers and Identify Bottlenecks: More hardware is not always the answer. Be vigilant for diminishing returns when adding resources. Use profiling and tracing tools (like Go execution tracing) to identify true bottlenecks (e.g., single-threaded WAL, GC pauses) rather than simply scaling up resources that aren't the limiting factor.
- Consider Go Runtime Tuning (Cautiously): If other optimizations don't resolve performance issues, tuning Go runtime garbage collection (
GOMEMLIMIT,GOGC) can provide significant benefits. However, this is an advanced technique that requires careful testing and a deep understanding of Go's memory management, as incorrect settings can lead to instability. - Plan for Horizontal Scaling Early: While vertical scaling can simplify initial deployments, plan for eventual horizontal scaling or sharding from the outset. Avoiding "too big to fail" moments by designing for distributed solutions will save significant headaches as systems grow. Also, consider offloading expensive historical queries to distributed backend systems like Grafana Cloud to reduce the load on local Prometheus instances.
Key Takeaways
- Vertical scaling has limits: Even with massive hardware, architectural and runtime bottlenecks (like WAL and Go GC) can prevent Prometheus from fully utilizing resources.
- Remote write requires careful tuning: Ignore
desired_shardsand calculate optimalshards,batch_size, andbuffer_sizebased on CPU and worst-case latency. - Network limits can be subtle: Be aware of single-flow egress bandwidth limits in cloud providers and consider disabling HTTP/2 if necessary to force multiple connections.
- High churn impacts compaction: Default Prometheus compaction can fail in dynamic environments; consider smaller, continuous 1-hour blocks to manage the write-ahead log effectively.
- Go runtime tuning is powerful:
GOMEMLIMITandGOGCcan mitigate garbage collection pauses, but apply these advanced optimizations cautiously after addressing other bottlenecks. - Simplicity has a cost: While "boring technology" is good, pushing a single instance to its absolute limits will eventually expose deep, complex performance challenges.
About the Speaker(s)
Chris Leavoy is an Observability Engineer at Etsy, based in Waterloo, Canada. He has been with Etsy for four years, contributing to the operation and scaling of their extensive monitoring infrastructure. His work involves tackling the challenges of maintaining robust observability for a global e-commerce marketplace, particularly at the scale discussed in this talk.
Bryan Boreham is a Distinguished Engineer at Grafana Labs and a long-standing Prometheus maintainer, having worked on the project's code for approximately seven years. At Grafana Labs, he focuses on scaling massively scalable storage solutions for metrics, logs, and traces, dealing with trillions of metric points and petabytes of logs. His deep expertise in Prometheus internals and Go runtime performance was instrumental in diagnosing and resolving the complex issues faced by Etsy.
Reviews
Dr. Zero (Offensive Security Researcher) — MUST SEE
This talk from Etsy and Grafana Labs isn't just another Prometheus scaling story; it's a brutal, honest, and deeply technical dive into the absolute limits of vertical scaling. They didn't just throw more hardware at the problem; they profiled Go runtime, uncovered subtle cloud provider network limitations, and re-engineered core Prometheus TSDB compaction and remote write logic. The findings are not just academic; they saved Etsy's Black Friday, led to upstream contributions, and provide actionable, low-level insights that challenge conventional wisdom. This is exactly the kind of deep, no-bullshit engineering exposé that advances the field.
Heather Calloway (CISO) — STRONG ACCEPT
This talk from Etsy and Grafana Labs provides a rigorous, unsentimental deep dive into the operational limits of Prometheus at extreme scale. While highly technical, it offers critical insights for any CISO or security leader concerned with the resilience of their core observability infrastructure. The detailed account of diagnosing and resolving performance bottlenecks, particularly those impacting data integrity and timely alerting during critical business periods like Black Friday, underscores the non-negotiable link between robust monitoring and institutional accountability for business continuity and incident response.