skip to content

How does Prometheus remote_write behave when the receiver is slow or down, and how does federation differ?

level: seniorimportance: must knowfreq 62%

answer

  1. One pushes, the other is scraped
  2. The sender reads its own recovery log
  3. Backlog costs memory, not correctness
  4. A missed pull is a permanent hole
  5. Truncation is where samples are finally lost

basics

~20 s

remote_write pushes samples out of the local write-ahead log, so a slow receiver grows the send queue and memory while scraping continues untouched; samples are lost only once the log is truncated past them. Federation instead pulls the latest value of selected series.

solid answer

~40 s

`remote_write` is a **push** run by the Prometheus process itself: a queue manager per endpoint reads the local WAL, shards samples across parallel senders, batches them and POSTs them to a URL, retrying with backoff. When the receiver slows, the queue grows, shard count climbs toward `max_shards`, memory climbs with it — but scraping and local storage are untouched, so the local server stays correct while the remote copy falls behind. Samples are only lost when the receiver stays down long enough that WAL truncation passes the unsent position. Federation is the opposite shape: a **pull**, where one Prometheus scrapes another's `/federate` endpoint with `match[]` selectors and `honor_labels: true`, receiving only the most recent value of each matched series at that instant. A missed federation scrape is simply a gap that never comes back.

code

yaml · 10 lines
yaml
remote_write:
  - url: https://metrics-lts.internal.example/api/v1/write
    queue_config:
      capacity: 10000
      max_samples_per_send: 2000
      min_shards: 4
      max_shards: 50
      batch_send_deadline: 5s
      min_backoff: 30ms
      max_backoff: 5s

go deeper

for a junior

Recall that one Prometheus can send its samples to an external store, and that a second Prometheus can also scrape a small selected set of series from a first one. Know that these are two different directions of travel.

for a middle

Explain the mechanics: remote_write is a push out of the write-ahead log with a sharded, batching, retrying queue; federation is an ordinary scrape of a /federate endpoint with match[] selectors that returns only the latest value of each series.

for a senior

Demonstrate failure-mode judgement. Say what degrades when a receiver is slow, name the lag and pending metrics you would look at, know where the durability boundary sits at log truncation, and be firm that federation is not a data-copying mechanism.

for a principal

Own the shape of the estate: which data leaves each server, what the sending servers cost in memory when the far end is unhealthy, and what standard you set so nobody solves a global-view problem with a federation job that silently loses data.

A single Prometheus is deliberately a single node with local storage. There are two supported ways to get its data somewhere else, and they behave nothing alike under failure. ## remote_write: a push driven from the log Configure a `remote_write` block and the Prometheus process starts a queue manager for that endpoint. That manager reads the local write-ahead log — the same log used for crash recovery — and streams what it finds to the configured `url` as compressed batches over HTTP. The mechanics that matter under load: - **Sharding.** Samples are distributed across parallel senders. The manager scales the shard count between `min_shards` and `max_shards` based on how far behind it is, so a receiver that is merely slow gets more concurrency thrown at it. - **Batching.** Each shard buffers up to `capacity` samples and sends up to `max_samples_per_send` at a time, flushing early when `batch_send_deadline` elapses. - **Retry with backoff.** Failed sends are retried with a delay growing between `min_backoff` and `max_backoff`. - **A write-time filter.** The block has its own filter stage so you can decline to ship a subset of series without touching what is stored locally. The crucial property is that **remote_write is not in the ingest path**. If the receiver is slow or unreachable, scraping continues, the local TSDB continues, queries against the local server continue to be correct. What degrades is the remote copy's freshness and the sending process's memory, because pending samples and additional shards are held in RAM. That leads directly to the durability bound. The queue reads from the WAL, and the WAL is truncated when blocks are cut and checkpoints are written. If the receiver stays down long enough that truncation passes the point the queue had reached, those samples are gone — they exist in local blocks, but they will never be shipped. The window is therefore roughly "how much WAL you keep", not "forever". The signals to watch are the ones the sending server exposes about itself: - `prometheus_remote_storage_highest_timestamp_in_seconds` minus `prometheus_remote_storage_queue_highest_sent_timestamp_seconds` is the lag in seconds — the single most useful number. - `prometheus_remote_storage_samples_pending` shows the backlog held in memory. - `prometheus_remote_storage_shards` against the configured maximum shows whether the manager has run out of concurrency to throw at the problem. - `prometheus_remote_storage_samples_failed_total` and the dropped counter separate "retrying" from "given up". ## Federation: a pull of the latest values Federation reuses the scrape machinery. A higher-level Prometheus is configured with an ordinary job whose `metrics_path` is `/federate` and whose `params` carry one or more `match[]` selectors; `honor_labels: true` keeps the labels the federated server returns rather than overwriting them with the scrape job's own. ```yaml scrape_configs: - job_name: federate-cellar-sites honor_labels: true metrics_path: /federate params: 'match[]': - '{__name__=~"site:.*"}' static_configs: - targets: ['prom-cellar-eu:9090', 'prom-cellar-us:9090'] ``` The response is one exposition-format payload containing the **most recent sample** of every series matching the selectors, evaluated at the moment of the request. Three consequences follow, and they are what the question is really testing: 1. **Resolution is lost.** Whatever the leaf server stored at fifteen-second granularity arrives at the parent's scrape interval. The intervening samples do not exist upstream and never will. 2. **A missed scrape is a permanent hole.** There is no queue, no retry of past data, no catch-up. If the parent could not reach the leaf for ten minutes, the parent simply has no data for those ten minutes. 3. **It is bounded by one HTTP request.** Matching a large series set means one enormous response the parent must parse inside a scrape timeout, and failing that timeout produces nothing at all rather than a partial result. ## Choosing between them, and the thing federation must never do | | remote_write | federation | remote_read | |---|---|---|---| | Direction | push, out of the sender | pull, into the parent | pull, on query | | Granularity | every sample | latest value per scrape | every sample | | Behaviour when the far end fails | queues, retries, then loses | permanent gap | query fails | | Suited to | long-term storage, a global store | a small aggregated roll-up | occasional access to older data | **Federation must never be used to copy an entire Prometheus's data to a global server.** That misuse is the single most common architectural mistake on this subject: it downsamples silently, it drops data whenever a scrape is missed, it collapses under the series volume, and it leaves the parent one timeout away from having nothing. The legitimate use is small and deliberate — pulling a handful of already-aggregated, cross-service series up to a higher-level view. Everything else, including long retention and a genuine global view, belongs on `remote_write`. `remote_read` completes the picture: it lets the local Prometheus fetch raw series from a remote store at query time and evaluate locally. It can be told to skip data the local server already holds and to fire only for queries carrying particular label matchers, because pulling wide ranges of raw series across the network at query time is exactly as expensive as it sounds.

  • A team wants a global view by federating every series from twelve regional Prometheus servers into one. Why is that a mistake?
    Federation returns only the latest value of each matched series per scrape, so raw resolution is destroyed on arrival, and any missed scrape is a permanent gap. Matching everything makes one enormous response that must be produced and parsed inside a scrape timeout, and the aggregate series volume then lands on a single node with the same limits as the ones it is aggregating. A global view belongs on remote_write into a store built for it.
  • Which numbers tell you a Prometheus remote_write queue is falling behind?
    The gap between `prometheus_remote_storage_highest_timestamp_in_seconds` and `prometheus_remote_storage_queue_highest_sent_timestamp_seconds` is the lag in seconds. Alongside it, `prometheus_remote_storage_samples_pending` shows the in-memory backlog, `prometheus_remote_storage_shards` against the configured maximum shows whether concurrency is exhausted, and the failed and dropped counters distinguish a receiver that is retrying from one that is rejecting.
  • If remote_write never blocks scraping, what is the failure that actually hurts the sending server?
    Memory. Pending samples and the extra shards spun up to chase a slow receiver are held in the process, so a long outage inflates resident memory on a server that also needs that memory for its head block. On a constrained node the sending Prometheus can be pushed into an out-of-memory kill by a problem in a completely separate system.

remote_write is a shop posting every receipt to head office from its own journal, catching up after the post strike ends; federation is head office phoning the shop each hour to ask for today's totals. Only one of them can catch up, and only one of them ever had every receipt.

saying these in an interview costs you the question

  • Thinks remote_write blocks scraping when the receiver is slow
  • Believes the remote_write backlog is buffered on disk forever
  • Uses federation to copy every series into a global server
  • Alerts on federated data as though it were raw and timely
  • Confuses a pull from /federate with a push over remote_write
  • Assumes remote_read is cheap for wide historical queries