Prometheus remote write at scale

Tune Prometheus remote write with verified queue settings, diagnose lag and retries, and calculate the spare throughput needed to recover after an outage.

Illustrative orange time-series curves on a dark grid with the Prometheus symbol
Editorial illustration, not a throughput benchmark or production dashboard.

Prometheus remote write reads samples from the write-ahead log and sends them through sharded queues to a compatible receiver. At scale, delivery depends on receiver capacity, sender resources and the time available to replay buffered data. Start with defaults, measure lag and failures, then test whether recovery throughput exceeds continuing ingestion.

What happens between scraping and remote storage?

Prometheus's local TSDB uses a write-ahead log (WAL) whether or not remote write is configured. Each remote-write destination reads the WAL into its own sharded sending queue. The Remote Write 1.0 specification defines Snappy-compressed protobuf messages sent by HTTP POST, with samples kept in timestamp order within each series.

  1. Scraped samples enter the local TSDB and its WAL.
  2. A destination's WAL reader feeds in-memory shard queues. Samples for a series are assigned to the same shard.
  3. Shards batch samples and send requests in parallel across different series.
  4. The receiver returns an HTTP status. The sender retries recoverable failures; permanent rejection requires investigation.

Prometheus calculates its own desired shard count from incoming samples, outstanding work and sending time. The receiver does not negotiate that count. A successful HTTP response reports a successful write under the protocol; replication, persistence and when data becomes queryable depend on the receiver. Verify those separately.

Buffering is finite. The official tuning guide warns that an endpoint outage beyond roughly two hours can lose unsent data when the server WAL is compacted. Treat this as a documented risk boundary, not a two-hour no-loss guarantee. Disk failure, permanent rejection or a configured sample-age limit can lose samples sooner. Local TSDB block retention is not a remote-write replay budget.

Which configuration should you start with?

Keep your existing scrape jobs. Merge this destination into the top-level remote_write list in prometheus.yml, replacing the example URL with the receiver's documented remote-write URL and adding its required authentication. The queue values below are the Prometheus 3.5.0 defaults. The explicit exception is retry_on_http_429: true, which opts into retrying rate-limited requests. That option is marked experimental in the pinned configuration reference.

yaml · prometheus.yml
remote_write:
  - name: primary
    url: https://metrics.example.com/api/v1/write
    protobuf_message: prometheus.WriteRequest
    remote_timeout: 30s
    queue_config:
      capacity: 10000
      min_shards: 1
      max_shards: 50
      max_samples_per_send: 2000
      batch_send_deadline: 5s
      min_backoff: 30ms
      max_backoff: 5s
      retry_on_http_429: true

The hostname is an example, not an xScaler endpoint. Use the URL, tenant identity and authentication issued for your deployment. Remote write, OTLP ingestion, remote read and the Prometheus query API are different interfaces; do not infer one URL from another. Retain TLS verification and keep credentials out of shared configuration examples.

Validate the merged file before a controlled reload. This command checks configuration, not network access, authentication or the receiver's ability to accept your data.

sh
promtool check config prometheus.yml

WAL compression is configured with the startup flag --storage.tsdb.wal-compression, not a storage.tsdb.wal_compression YAML field. It has been enabled by default since Prometheus 2.20.0. Compression reduces WAL storage use; it is not a repair or corruption-prevention mechanism. The storage documentation describes local filesystem requirements and recovery limitations.

How should you tune queues and shards?

Change one setting at a time against a measured bottleneck. Record sender CPU, memory, network use, series churn, request duration and receiver limits before increasing parallelism. A rate-limited receiver will not gain capacity because the sender opens more work.

Parameter meanings: Prometheus remote-write tuning and v3.5.0 configuration, checked 2 October 2026. These are decision criteria, not workload-specific recommendations.
SettingWhat it controlsWhen to reconsider it
capacityQueued samples per shard before WAL reading blocksShort request stalls fill queues; check memory first
max_shardsMaximum concurrent sending shardsSender is capped and receiver has demonstrated headroom
min_shardsInitial and minimum shard countMeasured startup lag matters before adaptive sharding catches up
max_samples_per_sendMaximum samples in a request batchMeasure request overhead and receiver payload limits
batch_send_deadlineFlush deadline for an unfilled batchTrade low-volume latency against request efficiency

When a shard queue fills, the WAL reader cannot continue feeding the destination's shards. More capacity can absorb a short stall, but cannot fix a sustained throughput deficit. Prometheus recommends capacity at three to ten times max_samples_per_send. The baseline's 10,000 divided by 2,000 is five.

Queue memory scales with shard count multiplied by capacity plus batch size, with additional memory for the series-label cache and other overhead. Do not turn that relationship into a byte estimate without measurement. High series churn can grow the cache even if the current active-series count looks acceptable. The upstream tuning guide explains these trade-offs.

Which metrics reveal lag and data loss?

Scrape the sending Prometheus instance's /metrics endpoint. Inspect each sender and remote_name independently so one healthy destination cannot conceal another's failure. The following names and behaviours are checked against the v3.5.0 queue manager. These examples cover ordinary float samples; deployments sending native histograms or exemplars must also monitor their corresponding counters.

Prometheus v3.5.0 metric definitions and send path, checked 2 October 2026. Verify names and labels on your sender before using the queries.
MetricUseful signalLimitation
prometheus_remote_storage_samples_pendingSamples pending in sender memoryDoes not count the entire unread WAL backlog
prometheus_remote_storage_samples_retried_totalSample retry activityThe same sample can contribute repeatedly
prometheus_remote_storage_samples_failed_totalPermanent or partial failures; hard-shutdown lossNot a counter of all transient errors
prometheus_remote_storage_samples_dropped_totalDrops labelled by reasonSeparate intentional relabelling from too_old and other loss
prometheus_remote_storage_shards_desiredCalculated sending concurrencyCompare with shards and shards_max; not proof of receiver capacity
prometheus_remote_storage_queue_highest_sent_timestamp_secondsTimestamp progress through the queueA maximum timestamp cannot establish complete delivery

Start with non-recoverable failures. This query retains sender and destination labels; do not sum it across every Prometheus instance before investigating.

promql
rate(prometheus_remote_storage_samples_failed_total{remote_name="primary"}[5m])

Inspect drops by reason. A rising dropped_series counter may reflect your deliberate write-relabel filter; too_old means age-based dropping. Unexpected reasons require investigation rather than a blanket allowance for dropped samples.

promql
rate(prometheus_remote_storage_samples_dropped_total{remote_name="primary"}[5m])

For a continuously active destination, the difference between wall-clock time and the queue's highest timestamp is a useful progress warning. Idle traffic, exporter timestamps and clock skew affect its interpretation. An absent series is not a zero-lag result; monitor the sender scrape separately.

promql
time() - prometheus_remote_storage_queue_highest_sent_timestamp_seconds{remote_name="primary"}

Avoid a universal alert such as pending > 10000. Capacity is per shard, while sample rate and the acceptable freshness delay vary by workload. Choose a sustained progress threshold from the investigation or alerting delay your team can tolerate, and test it during a staged outage. Keep an independent monitoring path if the remote backend is also where your alerts run.

What should you check when delivery stalls?

Use the HTTP response and sender logs to distinguish a retryable outage from data the receiver has rejected. The Remote Write 1.0 retry contract requires retrying HTTP 5xx with backoff, prohibits retrying ordinary HTTP 4xx and allows retries for HTTP 429. In the baseline configuration, 429 retries are explicitly enabled.

Diagnostic workflow based on the Remote Write 1.0 specification and Prometheus v3.5.0 sender behaviour, checked 2 October 2026. Observations narrow the investigation; they do not identify a cause by themselves.
ObservationCheck nextAction to test
429 with increasing retriesTenant rate limits and receiver capacityReduce offered work or arrange capacity; do not blindly add shards
5xx or connection failuresReceiver health, DNS, TLS and network pathRestore the failed dependency and watch recovery throughput
400 with failed samplesResponse body, labels and timestamp rejection detailsFix the rejected data or receiver contract; backoff alone cannot fix it
401 or 403Credentials, tenant identity and permissionsCorrect authentication; these are not ordinary retryable responses
Desired shards exceed the capSender resources and receiver headroomIncrease the cap only when the receiver can accept more
Fresh timestamp but missing seriesPermanent failures, relabelling and receiver rejectionsQuery expected series and counts across the affected window

Backoff spaces repeated attempts and doubles up to max_backoff. It does not guarantee that a fleet will avoid a recovery surge. Test restart and recovery behaviour across multiple senders if they share a receiver quota. Timestamp acceptance and out-of-order windows are receiver-specific; synchronise clocks and read the actual rejection reason instead of assuming a universal five-minute allowance.

How much headroom clears an outage backlog?

Use the rate of samples eligible for this destination after filtering. In a hypothetical workload generating 100,000 samples per second, a ten-minute delivery outage accumulates 60 million samples: 100,000 multiplied by 600 seconds. This assumes all of those samples remain replayable and the receiver accepts their timestamps.

After service returns, fresh samples still arrive. Net backlog-draining capacity equals sustainable successful delivery minus continuing ingestion. Recovery time equals backlog divided by that net capacity, provided it is positive. Use accepted, useful samples per second for the same workload, not HTTP attempt rate or a vendor's unrelated peak benchmark.

Illustrative constant-rate arithmetic, recomputed 2 October 2026: 60,000,000 / 25,000 = 2,400 seconds; 60,000,000 / 100,000 = 600 seconds. Times begin after delivery resumes. This is not measured xScaler or Prometheus throughput.
Successful deliveryFresh ingestionNet drainTime to clear
100,000 samples/s100,000 samples/s0 samples/sNever at these rates
125,000 samples/s100,000 samples/s25,000 samples/s40 minutes
200,000 samples/s100,000 samples/s100,000 samples/s10 minutes

Retries, resharding, receiver rejection and changing traffic can extend recovery or prevent it. A native-histogram-heavy workload is not interchangeable with the same count of float samples. Check the oldest replayable data and any sample_age_limit throughout recovery; a drain-time calculation cannot prove retention safety. The incident-readiness guide applies the same continuing-arrival constraint to a telemetry recovery exercise.

How do you preserve series identity?

Write relabelling filters or changes outgoing samples after external labels are applied. It does not remove the corresponding data from the local TSDB. The configuration reference distinguishes this from scrape-time metric relabelling.

For example, if a reviewed demo_debug_ metric family is not needed remotely, add this write_relabel_configs field to the same destination entry. It drops entire matching series from that destination. Do not create a second remote_write key or a duplicate primary destination when merging the example.

yaml · prometheus.yml
remote_write:
  - name: primary
    url: https://metrics.example.com/api/v1/write
    write_relabel_configs:
      - source_labels: [__name__]
        regex: "demo_debug_.*"
        action: drop

Dropping a label is different. Removing instance, pod or another distinguishing dimension can make previously distinct series share an identity; it does not add their values together. Preserve uniqueness and use explicit aggregation when that is the intended result. Review which queries and alerts rely on the detail before filtering, using the useful-labels guide.

High-availability replicas also need an explicit receiver contract. For example, Grafana Mimir's HA tracker uses configured cluster and replica labels to select a replica for ingestion. That is a Mimir feature with configuration requirements, not a remote-write protocol guarantee or evidence that every xScaler deployment enables it. Test failover and query continuity with the receiver you will use.

What should pass before a production rollout?

Run a bounded exercise against a staging receiver with representative labels, traffic and authentication. Preserve a known configuration to roll back to, and agree an acceptable freshness delay and loss budget before the test.

  1. Record normal ingestion, sender CPU and memory, shard count, failure counters and receiver query results. Include a continuously emitted canary series with known scrape timing.
  2. Introduce a short, controlled receiver outage while collection continues. Record the actual duration and response codes; do not delete or modify WAL files.
  3. Restore the receiver and measure successful recovery while fresh samples still arrive. Confirm the observed net drain can meet the recovery deadline.
  4. Compare raw canary timestamps, values and sample counts locally and remotely over a fixed window. A stepped query graph can conceal missing samples through lookback. Also check critical production-shaped series; one canary does not prove every series arrived.
  5. Exercise authentication rejection, rate limiting and any HA failover separately. Confirm the right alerts fire and that intentional filters do not hide unintended loss. Correcting a rejection does not automatically repair the rejected interval.
  6. Repeat the recovery check with a controlled sender restart and configuration reload. Verify the complete sample window again instead of assuming the replay behaviour is unchanged.
  7. Roll back if memory, receiver errors or data gaps exceed the agreed limits. Keep the configuration, versions, observations and query results with the rollout record.

A local Prometheus deployment can remain the simpler choice when its retention, query capacity and operational model meet your needs. Remote storage adds a network dependency and a second ingestion contract. If your required outage tolerance exceeds the sender's tested replay window, raising queue capacity is insufficient; evaluate a collection and buffering design that explicitly meets that requirement.

When evaluating xScaler, compare costs for the same retained series, resolution and query workload, then run these checks against the issued connection settings. Request explicit answers for receiver quotas, accepted timestamps, HA behaviour, data location and recovery support. Upstream Prometheus or Mimir documentation cannot establish those deployment-specific commitments.

Does a longer local retention period preserve unsent remote-write data?
Do not assume it does. Local TSDB block retention and the remote-write WAL replay path are different; test the replay window for your Prometheus mode and version instead of treating a retention setting as an outage guarantee.
Should I switch to Remote Write 2.0?
Only after confirming support for the protobuf message and required data types at the receiver. Prometheus 3.5.0 defaults to prometheus.WriteRequest; its configuration reference explicitly advises checking receiver compatibility before changing protobuf_message.
Can I use the same ingestion URL for OTLP and remote write?
Only if the provider explicitly documents that routing. They are different protocols, and the remote-write specification does not define a universal receiver path.
Will remote write send all my old TSDB blocks to a new destination?
It is not an automatic historical-block migration mechanism. Plan a separately supported backfill or migration process if you need historical data, and validate overlaps and receiver timestamp limits.

Technical sources checked on 2 October 2026, separately from the 10 May 2026 publication date. Pinned sources describe Prometheus 3.5.0; live receiver documentation is not an xScaler service commitment.