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.

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.
- Scraped samples enter the local TSDB and its WAL.
- A destination's WAL reader feeds in-memory shard queues. Samples for a series are assigned to the same shard.
- Shards batch samples and send requests in parallel across different series.
- 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.
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: trueThe 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.
promtool check config prometheus.ymlWAL 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.
| Setting | What it controls | When to reconsider it |
|---|---|---|
| capacity | Queued samples per shard before WAL reading blocks | Short request stalls fill queues; check memory first |
| max_shards | Maximum concurrent sending shards | Sender is capped and receiver has demonstrated headroom |
| min_shards | Initial and minimum shard count | Measured startup lag matters before adaptive sharding catches up |
| max_samples_per_send | Maximum samples in a request batch | Measure request overhead and receiver payload limits |
| batch_send_deadline | Flush deadline for an unfilled batch | Trade 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.
| Metric | Useful signal | Limitation |
|---|---|---|
| prometheus_remote_storage_samples_pending | Samples pending in sender memory | Does not count the entire unread WAL backlog |
| prometheus_remote_storage_samples_retried_total | Sample retry activity | The same sample can contribute repeatedly |
| prometheus_remote_storage_samples_failed_total | Permanent or partial failures; hard-shutdown loss | Not a counter of all transient errors |
| prometheus_remote_storage_samples_dropped_total | Drops labelled by reason | Separate intentional relabelling from too_old and other loss |
| prometheus_remote_storage_shards_desired | Calculated sending concurrency | Compare with shards and shards_max; not proof of receiver capacity |
| prometheus_remote_storage_queue_highest_sent_timestamp_seconds | Timestamp progress through the queue | A 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.
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.
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.
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.
| Observation | Check next | Action to test |
|---|---|---|
| 429 with increasing retries | Tenant rate limits and receiver capacity | Reduce offered work or arrange capacity; do not blindly add shards |
| 5xx or connection failures | Receiver health, DNS, TLS and network path | Restore the failed dependency and watch recovery throughput |
| 400 with failed samples | Response body, labels and timestamp rejection details | Fix the rejected data or receiver contract; backoff alone cannot fix it |
| 401 or 403 | Credentials, tenant identity and permissions | Correct authentication; these are not ordinary retryable responses |
| Desired shards exceed the cap | Sender resources and receiver headroom | Increase the cap only when the receiver can accept more |
| Fresh timestamp but missing series | Permanent failures, relabelling and receiver rejections | Query 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.
| Successful delivery | Fresh ingestion | Net drain | Time to clear |
|---|---|---|---|
| 100,000 samples/s | 100,000 samples/s | 0 samples/s | Never at these rates |
| 125,000 samples/s | 100,000 samples/s | 25,000 samples/s | 40 minutes |
| 200,000 samples/s | 100,000 samples/s | 100,000 samples/s | 10 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.
remote_write:
- name: primary
url: https://metrics.example.com/api/v1/write
write_relabel_configs:
- source_labels: [__name__]
regex: "demo_debug_.*"
action: dropDropping 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.
- 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.
- Introduce a short, controlled receiver outage while collection continues. Record the actual duration and response codes; do not delete or modify WAL files.
- Restore the receiver and measure successful recovery while fresh samples still arrive. Confirm the observed net drain can meet the recovery deadline.
- 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.
- 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.
- 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.
- 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.
- Prometheus: Remote Write 1.0 specification 2 Oct 2026
- Prometheus: remote-write tuning 2 Oct 2026
- Prometheus 3.5.0: configuration defaults 2 Oct 2026
- Prometheus 3.5.0: configuration reference 2 Oct 2026
- Prometheus 3.5.0: queue metrics and sending implementation 2 Oct 2026
- Prometheus 3.5.0: local storage and remote integrations 2 Oct 2026
- Grafana Mimir: configuring HA deduplication 2 Oct 2026