Skip to content
SRE & DevOpsDeep Dive Published Updated 10 min readViews unavailable

Prometheus Remote Write: Backpressure, WAL Recovery, and Delivery

Operate Prometheus remote write by sizing queues, detecting backlog, handling WAL replay and receiver outages, and validating protocol compatibility.

Prometheus remote write forwards locally ingested samples to a remote-compatible storage system. It is useful for centralized retention and multi-cluster querying, but it introduces another asynchronous boundary: Prometheus can scrape successfully while the remote receiver is slow, overloaded, unreachable, or rejecting the data format. A green scrape target does not prove that its samples have reached the remote backend.

Reliable operation requires understanding the queue between Prometheus’ write-ahead log (WAL) and each remote endpoint, measuring how quickly a backlog is growing, and knowing which failures are retryable. Remote write is not an indefinitely durable message broker, and an HTTP success response does not necessarily prove that the receiver has committed data to durable storage. Confirm the sender, receiver, protocol version, and persistence contract as a pair.

Follow the sample through the pipeline

The local TSDB and the remote destination are separate storage paths. Prometheus scrapes targets and appends samples locally; remote write reads the WAL, builds an in-memory queue divided into shards, and sends batches to the configured endpoint. That local-first design helps keep scraping independent of momentary backend latency, but it does not make the remote queue unbounded or the receiver durable.

exporters -> scrape -> local TSDB / WAL -> remote-write queues -> receiver -> remote storage
                 |                           |
                 +-- local queries            +-- asynchronous delivery and retry

Each configured destination has a queue. Prometheus dynamically adjusts shard concurrency according to incoming sample rate, outstanding work, and send latency. If a shard falls behind and its in-memory capacity fills, reading from the WAL can block for that destination. More shards are not automatically better: they increase concurrent requests and memory use, and can overwhelm a rate-limited receiver.

Remote write does not forward arbitrary application pushes into Prometheus. The standard model is that a Prometheus-compatible sender scrapes or otherwise ingests metrics and forwards them. Local retention and remote retention are also independent: deleting local blocks does not necessarily delete remote samples, and remote-write filters do not erase the local series.

Treat backlog as a time budget

The most actionable built-in backlog signal is prometheus_remote_storage_samples_pending. Alert on its rate of change and duration, not only a static count. If the sender ingests 50,000 samples per second and the queue grows by 10,000 per second, the receiver is accepting less than the offered rate; a growing backlog is an outage trend even before samples are lost. Estimate the remaining recovery window from observed WAL behavior, disk headroom, and actual send throughput.

Prometheus’ current remote-write tuning guidance states that failed requests are retried, but samples not sent may be lost if an endpoint remains down for more than two hours and WAL segments are compacted. Treat that as an important design boundary, not as a guaranteed two-hour durable queue: validate the exact version and storage behavior, and do not use it as the sole buffer for a destination that may be unavailable longer. If the remote system is the only source of long-term retention, establish a separate failure strategy and explicitly define the acceptable data-loss window.

When an endpoint recovers, it must drain the backlog faster than new samples arrive. If the receiver can accept 60,000 samples per second while Prometheus continues ingesting 50,000, only 10,000 samples per second of net catch-up remain. Recovery time can therefore exceed outage duration substantially. Forecast receiver capacity, rate limits, network bandwidth, sender CPU, and available WAL/disk space together. A receiver that returns from outage at normal steady-state capacity may never clear the backlog if it cannot temporarily exceed the current ingestion rate.

Tune the queue as a bounded system

Queue settings live under each remote_write entry’s queue_config. Common controls include capacity per shard, max_shards, min_shards, max_samples_per_send, and retry backoff. Prometheus documentation recommends a capacity around three to ten times max_samples_per_send as a starting relationship. That is not a workload-specific memory target; series churn, actual batch size, number of queues, and active shards all matter.

remote_write:
  - name: long-term-store
    url: https://metrics.example.net/api/v1/write
    queue_config:
      capacity: 10000
      max_shards: 50
      min_shards: 1
      max_samples_per_send: 2000
      batch_send_deadline: 5s
      min_backoff: 30ms
      max_backoff: 5s

These illustrative values match common Prometheus defaults in current documentation; they are not a prescription for a receiver or production traffic profile. Validate the actual config against the deployed release with promtool check config prometheus.yml. A value that helps an isolated benchmark can increase memory pressure or overload the real receiver during a backlog recovery.

Tune in a controlled order:

  1. Confirm the receiver, network path, authentication, rate limits, and remote protocol are healthy. Queue tuning cannot correct an invalid URL, expired credentials, or a backend rejecting samples.
  2. Measure ingestion rate, pending samples, request latency and failures, sender CPU/network, Prometheus memory, and receiver saturation during normal load and a controlled slowdown.
  3. Adjust batch size and the number of concurrent shards within the receiver’s supported request size and ingestion budget. Increasing max_samples_per_send may reduce request overhead but can increase latency and exceed receiver limits.
  4. Adjust queue capacity only after calculating the memory cost per active shard and queue. Larger capacity absorbs short delays but can consume substantial memory and take longer to drain after resharding.
  5. Test a receiver outage and recovery. Verify that queue metrics return to baseline, local scraping remains healthy, and receiver capacity drains the backlog without causing a retry storm.

Prometheus automatically chooses shard count within configured bounds. Increase max_shards only when measurements show that more parallelism can increase throughput and the receiver has spare capacity. Decrease it when a backend is being overwhelmed or memory is constrained. Raising min_shards may help a known high-volume destination start at needed parallelism, but it also commits resources from startup. Prefer evidence from the pending-sample trend over guessing.

Understand failure and retry semantics

The Remote Write 1.0 specification defines HTTP-based batched requests and sender retry behavior. In general, transient server-side failures (5xx) should be retried with backoff; invalid or permanently rejected requests (most 4xx) should not be retried indefinitely. HTTP 429 is special: a receiver may use it for overload, and retry behavior depends on sender version and configuration. Current Prometheus configuration documents an experimental retry_on_http_429 option that defaults to false. Confirm the exact response code, Retry-After behavior, sender configuration, and receiver contract before changing it.

Do not classify all errors as “network, retry.” A persistent 401 or 403 usually requires correcting credentials or permissions; 400 can indicate invalid samples or limits; 413 may indicate request-size limits; 429 may be backpressure; 5xx may indicate a receiver outage or overload. Logs and metrics should retain the endpoint name and status category without exposing secret headers.

The Remote Write 2.0 specification adds a new protobuf message, written-count response headers, and protocol improvements, but its current upstream spec is still marked experimental. Prometheus exposes a protobuf_message setting, and its configuration docs explicitly advise consulting or testing the receiver before changing the payload. Do not switch a fleet to the v2 message solely because the sender accepts the config: first verify receiver support, partial-write behavior, headers, histograms/exemplars, rollback path, and how a mixed-version rollout negotiates or fails.

An HTTP 2xx is evidence of protocol-level acceptance, not proof of permanent storage or query visibility. In the Remote Write 2.0 spec, “written” means the receiver accepted the data; whether it has persisted it is receiver-defined. Test the receiver’s durability and query-latency guarantees separately. If receiver retries can replay a request after an ambiguous connection failure, validate how duplicate samples are handled by that backend rather than assuming exactly-once delivery end to end.

Preserve series identity and control fan-out

Remote-write labels form the identity of each time series. Configure stable external labels such as cluster or environment identity deliberately, and avoid assigning the same label set to independent Prometheus senders unless the receiver has a documented high-availability deduplication strategy. Two HA replicas may scrape the same targets and send overlapping samples; deduplication is generally backend-specific and is not guaranteed by the transport protocol.

Each additional remote destination creates another queue and network workload. If one receiver exists for short-term regional analytics and another for long-term storage, budget for two independent queues, their retry behavior, credentials, rate limits, and traffic costs. A slow endpoint must not be assumed to have no effect just because the other endpoint is healthy; monitor each named destination separately.

Use write_relabel_configs only when the intended result is to change what one remote destination receives. These rules are applied after external labels and do not remove matching series from the local TSDB. Keep the local and remote query contracts explicit: a local dashboard may still show a metric that is intentionally filtered from remote storage. For large fleets, review filtering and cardinality as both correctness and cost controls.

Observe the sender and receiver together

At minimum, alert on sustained pending samples, remote write failures, oldest-unsent or queue progress where available, and loss/retry signals exposed by the deployed version. Alert on a remote-write destination that has stopped advancing even when up is healthy for every scrape target. Pair sender telemetry with receiver request rate, accepted and rejected samples, ingestion latency, throttling, storage health, and query visibility.

At incident start, preserve the timeline and answer these questions:

  • Are samples still being scraped and appended locally, or has local ingestion also failed?
  • Is pending work rising, flat, or draining for each destination?
  • Which HTTP status, timeout, TLS, DNS, authentication, or request-size error is reported?
  • Can the receiver accept more than current ingestion to catch up, and does it have memory and disk headroom?
  • Is Prometheus under CPU, memory, network, or WAL disk pressure because of retries or series churn?
  • Are samples absent from the remote query path, or merely delayed by receiver ingestion/indexing?

Avoid restarting Prometheus reflexively: a restart can interrupt in-flight requests, trigger replay, and hide useful queue evidence. Avoid increasing shard limits during overload until receiver capacity is understood. Preserve logs and time-series evidence, validate a minimal remote-write request where safe, and apply one bounded change at a time.

Production acceptance checklist

  • Local scrape health and remote-delivery health have separate dashboards and alerts for every named destination.
  • The accepted data-loss window is explicit and does not rely on WAL replay as an indefinite durable queue.
  • Queue capacity, shards, batch size, and retry behavior are measured against sender memory and receiver throughput.
  • Outage and catch-up tests prove that the receiver can drain backlog faster than new ingestion without overload.
  • HTTP 401/403, 400/413, 429, 5xx, DNS, TLS, and timeout paths have distinct runbook actions.
  • Sender and receiver protocol versions, message format, partial writes, and native histogram/exemplar support are documented and tested.
  • Series labels identify cluster and HA replicas consistently, with any deduplication verified as a backend feature.
  • Credentials use the supported authentication mechanism and are not embedded as cleartext in a publicly accessible repository.
  • promtool check config passes, and local-versus-remote filtering behavior is tested before production rollout.

Remote write is an asynchronous forwarding pipeline with bounded queues, receiver-specific semantics, and a finite recovery horizon. Production reliability comes from measuring its backlog and failure budget, sizing both ends together, validating delivery behavior under outage, and keeping local health distinct from remote acceptance.

Related:

Sources:

Comments