A single Prometheus server goes a long way. It scrapes targets, stores samples in a local time-series database and answers PromQL, with no cluster to operate. Eventually one of four things breaks: the server runs out of memory because there are too many active series, one server cannot see every cluster, queries over months of data are too slow or the data has expired, or losing that server means losing alerts. Scaling Prometheus means deciding which of these you actually have and applying the cheapest fix that solves it.

This article is that decision path, in the order you should take it: measure, reduce, shard, pair for availability, then build a global view with a horizontally scaled backend. Fundamentals of the scrape model, PromQL and the local TSDB are in Prometheus Deep Dive; this page assumes them and focuses on scale. Product details were checked against the Prometheus, Grafana Mimir and Thanos documentation on 2026-10-02.

Advertisement

What actually limits one server

Prometheus keeps the most recent data, roughly the last two hours, in an in-memory head block backed by a write-ahead log, and writes older data to immutable blocks on disk. Memory is dominated by the number of active series in the head, plus the series that appeared and disappeared during the head window. That second part is churn: every pod restart creates new series because the pod label changes, and the old ones linger until the head is compacted. A cluster that redeploys often can carry several times its steady-state series count.

Disk and CPU rarely bind first. Sample ingestion is cheap; compression stores most samples in one or two bytes. The usual order of failure is memory from series and churn, then slow queries that touch many series, then retention limits from local disk, then the lack of a global view. Do not guess which: measure.

# Active series in the head block, and how fast they churn
prometheus_tsdb_head_series
rate(prometheus_tsdb_head_series_created_total[1h])

# Samples ingested per second
rate(prometheus_tsdb_head_samples_appended_total[5m])

# Which jobs bring the most series (post-relabel samples per scrape)
topk(10, sum by (job) (scrape_samples_post_metric_relabeling))

# Which metric names have the most series (expensive: run ad hoc, not on a dashboard)
topk(20, count by (__name__) ({__name__=~".+"}))

# Slowest rule groups: evaluation time against interval
max by (rule_group) (prometheus_rule_group_last_duration_seconds)

Record these as a dashboard per server and watch them over weeks. The ratio of series created per hour to active series tells you how much of your memory is churn. The top jobs and metric names tell you where to cut. For the reasons cardinality grows and how to budget it, see Metric cardinality.

Step 1: reduce before you scale

Every scaling architecture below costs more per series than the server you have. Removing series is the cheapest scaling step and usually the largest. Three sources dominate in practice: histograms with many buckets multiplied by many label combinations, labels with unbounded values (user IDs, request IDs, full URLs), and exporters that expose hundreds of metrics nobody queries.

Use metric_relabel_configs to drop unused metrics and strip unbounded labels at scrape time. Add per-job limits as guards: sample_limit fails the whole scrape if a target exposes more samples than allowed after relabelling, label_limit fails it if any sample has too many labels, and target_limit marks targets failed when service discovery finds more than the limit. A failed scrape is visible through the up metric, which is the point: a team that ships a label explosion finds out at once instead of taking down shared monitoring.

scrape_configs:
  - job_name: kubernetes-pods
    sample_limit: 50000        # whole scrape fails above this: a guard, not a quota
    label_limit: 30
    target_limit: 2000
    kubernetes_sd_configs: [{ role: pod }]
    metric_relabel_configs:
      # Drop a histogram nobody queries; it was 22% of series
      - source_labels: [__name__]
        regex: http_server_request_size_bytes_bucket
        action: drop
      # Strip an unbounded label before ingestion
      - regex: request_id
        action: labeldrop

Recording rules help the query side. A dashboard that sums a high-cardinality histogram every 30 seconds over seven days is expensive; a recording rule that pre-aggregates to the labels the dashboard uses turns it into a cheap lookup. Recording rules do not reduce ingestion, because the raw series are still stored, unless you also stop shipping the raw series downstream.

Advertisement

Step 2: shard the scraping

When one server cannot hold the series after reduction, split the targets. Functional sharding gives each team, cluster or domain its own server: infrastructure, payments, search. It is simple, and failures stay local, but shards are uneven and queries across them need a global layer. Hashmod sharding splits one large job across N servers by hashing a target label, usually the address, and keeping only the targets whose hash matches the server's index.

# Prometheus shard N of 4 scrapes only the targets whose address hashes to N
global:
  external_labels:
    cluster: prod-eu-1
    __replica__: replica-a      # Mimir HA tracker default label; Thanos usually uses "replica"

scrape_configs:
  - job_name: big-job
    kubernetes_sd_configs: [{ role: pod }]
    relabel_configs:
      - source_labels: [__address__]
        modulus: 4
        target_label: __tmp_hash
        action: hashmod
      - source_labels: [__tmp_hash]
        regex: "2"              # this server is shard 2
        action: keep

Hashmod shards are balanced by target count, not series count, so one heavy target can still skew a shard. Changing N reshuffles most targets, which resets their series and briefly multiplies churn; grow in planned steps. Alerting rules on a single shard see only that shard's targets, so any rule that aggregates across the job must run in the global layer.

Step 3: pairs for availability

A single server is a single point of failure for alerting. The standard answer is an HA pair: two identical servers scraping the same targets, each sending alerts to a clustered Alertmanager, which deduplicates notifications. Each server gets an external label identifying the replica. The data from the two replicas differs slightly, because scrapes happen at slightly different times, so whatever reads both must deduplicate.

Thanos does this at query time: the querier is started with a replica label, and series that differ only by that label are merged. Grafana Mimir does it at write time: the distributor's HA tracker elects one replica per cluster as leader and accepts samples only from it. By default it identifies clusters and replicas by the cluster and __replica__ labels, and if the leader sends nothing for 30 seconds it switches to the other replica, so at a 15-second scrape interval a failover typically loses about one scrape. From Mimir 3.0, memberlist is the recommended store for the election state.

Step 4: the global view

Once data lives on more than one server, something has to answer queries across all of them and keep data longer than local disks allow. There are three families of answer.

ApproachHow it worksGood forLimits
FederationA higher-level Prometheus scrapes /federate on lower ones, usually pulling only aggregated seriesA small global dashboard of recording-rule outputsNot for raw data; the top server becomes the bottleneck and a gap source
Thanos sidecar and storeA sidecar uploads each local block to object storage and serves recent data; a querier fans out to sidecars and store gatewaysKeeping existing servers; long retention in object storageQueries depend on every sidecar being reachable; recent data stays on the edge
Remote write backendServers push samples to Mimir, Cortex or Thanos Receive, which replicate, store in object storage and serve queries centrallyMany clusters, multi-tenancy, central limitsA large distributed system to operate; network and backend become part of the ingest path
Targetscluster ATargetscluster BPrometheus A1, A2HA pair, or --agentPrometheus B1, B2hashmod shardsscrapescrapeDistributorHA dedup, limitsremote_writeremote_writeIngestersor Kafka in 3.0Object storageTSDB blocksflush blocksQuerierPromQL engineStore gatewayreads blocksrecent dataQuery frontendsplit, cache, queueGrafana, rulesdashboards, alertsLocal Prometheus servers scrape; a horizontally scaled backend stores, deduplicates and queries the global view
Remote write to a Mimir-style backend: local servers scrape, the distributor deduplicates and enforces limits, and queries hit a split read path.

Once you have a remote write backend, the edge servers no longer need local queries or long retention. Prometheus can then run in agent mode, started with the --agent flag in Prometheus 3, which keeps only discovery, scraping and remote write with a write-ahead log. It uses far less memory than a full server and cannot answer queries or evaluate rules locally, so rules move to the backend's ruler.

Remote write: tuning and watching the pipe

Remote write reads from the write-ahead log and sends batches through a set of parallel shards, each with its own queue. Prometheus scales the shard count between the configured minimum and maximum based on how far behind it is. If the backend is slow or down, samples wait in the WAL; if the outage outlasts what the WAL holds, roughly the head window of about two hours, data is lost.

remote_write:
  - url: https://mimir.example.internal/api/v1/push
    headers: { X-Scope-OrgID: payments }
    queue_config:               # example values for a busy server; tune from the metrics below
      capacity: 10000
      max_shards: 50
      min_shards: 1
      max_samples_per_send: 2000
      batch_send_deadline: 5s
      min_backoff: 30ms
      max_backoff: 5s
    write_relabel_configs:
      - source_labels: [__name__]
        regex: "go_gc_.*|process_.*"
        action: drop            # keep locally for debugging, do not pay to ship it

The values above are examples for a busy server, not defaults; check the documentation for your version's defaults and change one setting at a time. Larger batches mean fewer requests and more memory per shard. More shards mean more parallelism and more connections to the backend. Use write_relabel_configs to keep debugging metrics local rather than paying to store them centrally. Then alert on the pipe itself:

# Remote write falling behind: seconds between newest sample and newest sample sent
(
  max_over_time(prometheus_remote_storage_highest_timestamp_in_seconds[5m])
- ignoring(remote_name, url) group_right
  max_over_time(prometheus_remote_storage_queue_highest_sent_timestamp_seconds[5m])
) > 120

# Shards pinned at the maximum: the queue cannot keep up
prometheus_remote_storage_shards >= prometheus_remote_storage_shards_max

Inside a horizontally scaled backend

Mimir, which grew out of Cortex, splits the work into components that scale independently. Distributors validate incoming samples, apply per-tenant limits such as ingestion rate and maximum series, run the HA tracker, and hash each series to a set of ingesters. Ingesters hold recent data in memory, much like a Prometheus head block, and periodically upload TSDB blocks to object storage. The compactor merges and deduplicates blocks. Store gateways serve queries over old blocks, queriers evaluate PromQL across ingesters and store gateways, and the query frontend splits long queries by time, caches results and queues work fairly between tenants.

Mimir 3.0, released in November 2025, added an ingest storage architecture that puts Kafka between the write and read paths. Distributors write to Kafka and ingesters consume from it, so a burst of heavy queries no longer slows ingestion and an ingester failure is less likely to break reads. The same release made the streaming Mimir Query Engine the default, which Grafana reports cuts peak query memory substantially compared with the Prometheus engine. Treat the Kafka layer as a real dependency with its own capacity planning; the classic architecture without it remains an option.

Thanos reaches a similar shape from a different starting point. Its compactor also downsamples old blocks to five-minute and one-hour resolution, which keeps long-range queries fast; how to do that without misleading readers is covered in metric downsampling.

Worked example: from one server to forty million series

Year one. One Prometheus pair per Kubernetes cluster, three clusters, about 1.5 million active series each. A dashboard showed churn adding 60% on deploy days. The team dropped two unused histograms and a pod_template_hash label, cutting series by 35%, and added sample_limit to every job. Memory headroom returned without new infrastructure.

Year two. Twelve clusters and a requirement for 13 months of retention and cross-cluster SLO dashboards. Federation could not carry raw data and sidecars in twelve regions made every global query depend on twelve networks. The team deployed Mimir centrally, one tenant per business unit, and switched edge servers to remote write with the HA tracker. Per-tenant series limits replaced the honour system. SLO burn-rate rules, described in burn-rate alerting, moved to the Mimir ruler because they span clusters; node-level alerts stayed at the edge so they still fire if the central system is unreachable.

Year three. Forty million series. Edge servers moved to agent mode, saving most of their memory. Query load from a capacity-planning tool began to slow ingestion during business hours, the problem ingest storage targets, so the team planned the move to the Kafka-based architecture with load tests against a copy of production traffic.

Failure modes

  • Cardinality explosion. One deploy adds a label with unbounded values and doubles series. Guard with sample_limit at the edge and per-tenant series limits centrally, and alert on series growth rate, not just totals.
  • Remote write backlog. The backend slows, shards max out, the WAL fills, and data is lost after the head window. Alert on lag and shard saturation, and size the backend for peak, not average.
  • Duplicate or missing data from HA pairs. Replica labels differ between pairs or are missing. Standardise external labels in one shared config.
  • Global alerting as single point of failure. All alerts run centrally, and the central system's outage silences them. Keep critical, local alerts at the edge.
  • Expensive queries. A dashboard regex over all series exhausts queriers. Use the query frontend's limits and recording rules for hot dashboards.
  • Resharding churn. Changing the hashmod modulus resets most series. Change it rarely and in planned steps.

What to do next

  1. Build the measurement dashboard: head series, series created per hour, samples per second, top jobs and rule duration.
  2. Drop unused metrics and unbounded labels, then set sample_limit, label_limit and target_limit on every job.
  3. Run every alerting Prometheus as an HA pair with a replica external label and a clustered Alertmanager.
  4. Shard functionally first; use hashmod only for a single job that is too large for one server.
  5. When you need a global view or long retention, choose between Thanos sidecars and a remote write backend using the table above.
  6. If you adopt remote write, alert on lag and shard saturation, and consider agent mode at the edge.
  7. Keep SLO rules, described in SLIs, SLOs and error budgets, in the global layer and node-level alerts at the edge.
Key takeaway: Scale Prometheus in order. First measure head series, churn and top jobs, and remove series before adding infrastructure. Guard every job with sample, label and target limits. Shard by function before hashmod, and run alerting servers as HA pairs. Thanos deduplicates replicas at query time and Mimir's HA tracker does it at write time. For a global view and long retention, use Thanos sidecars or remote write to a backend such as Mimir. Mimir 3.0 adds Kafka-based ingest storage that separates reads from writes. Alert on remote write lag, and keep critical alerts at the edge.