Scaling Prometheus & Long-Term Storage

A single Prometheus server is remarkably capable: it can ingest millions of active series on one machine, with fast local queries. But it's designed as a single-node, local-storage system. Retention is limited by local disk, there's no built-in clustering or replication, each server sees only the targets it scrapes, and one node's memory caps how many series it can handle.

As organizations grow, they need months or years of history, high availability, and a global view across clusters and regions. The ecosystem answers with remote write into horizontally scalable, object-storage-backed systems (Thanos, Grafana Mimir, VictoriaMetrics, managed services), plus disciplined cardinality management, since series count drives both cost and stability.

TL;DR

Quick Example

Prometheus in each Kubernetes cluster forwarding to a central Mimir:

Or run Prometheus in agent mode (--agent), which scrapes and forwards without local querying, rule evaluation, or long-term storage, using less memory.

Core Concepts

The Local TSDB

Prometheus stores recent samples in memory (the head block) with a write-ahead log, then compacts them into immutable 2-hour blocks on disk, merged into larger blocks over time. Compression is excellent (around 1–2 bytes per sample). Key settings:

Local storage isn't replicated or clustered. Losing the disk loses the history.

High Availability

Prometheus HA is simple: run two (or more) identical replicas with different replica external labels. Both scrape the same targets and evaluate the same rules. Alertmanager deduplicates their alerts, and query layers (Thanos Query, Mimir's HA tracker, Grafana data sources) deduplicate or pick one replica's data. Small gaps from scrape timing differences are acceptable.

Federation

A higher-level Prometheus scrapes the /federate endpoint of lower-level servers for selected, usually pre-aggregated series (recording rule outputs such as job:http_requests:rate5m). It suits hierarchical rollups, such as per-datacenter to global. It isn't a way to copy all raw data centrally: federating everything overwhelms the parent.

Remote Write and Remote Read

remote_write streams every ingested sample (after optional write_relabel_configs filtering) to an external endpoint in near real time, with a disk-backed WAL buffer and sharded queues. It's now the standard way to feed long-term and global storage. Remote write 2.0 improves efficiency and metadata. Many backends also accept OTLP metrics directly.

Long-Term Storage Options

All provide durable object-storage-backed retention (S3, GCS, Azure Blob), horizontal scaling, and PromQL-compatible query APIs for Grafana.

Downsampling and Retention Tiers

Keeping raw 30-second resolution for a year is expensive and slow to query. Thanos Compactor downsamples to 5-minute and 1-hour resolution for older data, and other systems offer tiered retention per tenant or metric. Match retention to use: high resolution for weeks (debugging), low resolution for years (capacity planning, trends).

Cardinality: The Real Scaling Limit

Resource use and cost are driven by active series, not by the number of metric names or targets. Cardinality explosions usually come from a new label with unbounded values (user IDs, request paths, pod names in a fast-churning deployment).

Measure:

Also check the TSDB status page (/tsdb-status), mimirtool analyze, and Grafana's cardinality management dashboards.

Control:

Best Practices

Keep Prometheus Close to Targets

Run Prometheus (or agents) per cluster or region, scraping locally, and ship data centrally with remote write. Scraping across regions adds latency, cost, and failure modes.

Evaluate Critical Alerts Locally

Alert rules evaluated in the local Prometheus keep working even if the central store or network is down. Use the global store for cross-cluster alerts and long-range queries.

Use Consistent External Labels

cluster, region, env, and replica external labels make global queries, deduplication, and multi-tenant routing possible. Define them from day one.

Budget Series Like Money

Track series per team or service, include cardinality in code review for instrumentation changes, and alert on sudden series growth. In managed services, series and samples translate directly into cost. See cloud costs.

Common Mistakes

Scaling Up One Giant Prometheus Forever

Adding RAM to a single server handling everything eventually hits limits: slow restarts (WAL replay), huge queries, and a single point of failure. Shard by cluster or function, and centralize via remote write.

Federating Raw Data

Federating all series into a central Prometheus recreates the scaling problem in one place, with extra latency. Federate aggregates only, or use remote write with a scalable backend.

Ignoring Remote Write Backpressure

If the backend is slow or down, remote write queues grow and the WAL buffer eventually drops data. Monitor prometheus_remote_storage_samples_pending, failed samples, and shard counts, and alert on remote write lag.

FAQ

How long can Prometheus retain data?

As long as your disk allows, but local storage isn't replicated, and huge retention makes queries and restarts slow. For months or years of data, use remote write to a long-term store with object storage, keeping local retention to days or weeks.

Thanos, Mimir, or VictoriaMetrics?

Thanos fits well if you already run Prometheus servers and want to add global querying and object-storage retention incrementally via sidecars. Mimir suits large, multi-tenant, remote-write-centric platforms. VictoriaMetrics emphasizes resource efficiency and operational simplicity. Managed services remove the operational burden entirely. All support PromQL and Grafana.

Is Prometheus agent mode a replacement for Prometheus?

For edge or per-cluster collection that only forwards data, yes: agent mode scrapes and remote-writes with a smaller footprint. It doesn't store data locally for queries or evaluate rules, so you need a backend that handles querying and alerting.

What is "active series" and why does it matter?

An active series is a unique metric-plus-label combination that has received samples recently, held in memory. It determines memory usage, ingestion cost, and query performance far more than sample count. Keeping it under control is the central scaling concern.

Related Topics

References