Designing Prometheus Remote Write for Long-Term Storage: An Architecture Note
An architecture note on Prometheus Remote Write: requirements, minimal design, trust boundaries, the metrics that prove it works, failure modes, and when to pivot to Remote Read.
01 Jul 2025, 23:10 UTC

The problem Remote Write solves
Prometheus stores data in a local time-series database (TSDB) sized for days or weeks, not years. When you need multi-year retention, cross-cluster querying, or a copy of metrics outside the cluster that produced them, the supported answer is Remote Write: Prometheus pushes every sample it scrapes to an external system such as Thanos, Cortex, Mimir, or a managed metrics service. The takeaway for planning: Remote Write is a best-effort, one-way export pipeline, not a backup or replication mechanism, and your design should treat it that way from the start.
Requirements to pin down first
Before writing configuration, answer three questions:
- Retention: How long must data live beyond local retention? This determines whether the remote store needs object storage backing (Thanos/Mimir style) or a simpler receiver is enough.
- Query pattern: Will dashboards and alerts keep hitting the local Prometheus, or will users query the long-term store directly? Heavy remote querying changes the design (see the pivot section below).
- Failure tolerance: Is it acceptable to lose samples during extended remote-store outages? Remote Write buffers in the local write-ahead log (WAL), but that buffer is finite.
Also budget for the overhead: Remote Write adds CPU cost for Snappy compression and Protobuf serialization, plus continuous network egress roughly proportional to your ingestion rate. On a server scraping hundreds of thousands of active series, this is measurable and should appear in capacity planning.
The smallest suitable design
The minimal architecture is one Prometheus server, one remote_write endpoint, and local retention kept short (for example 15 days) once you trust the remote copy. Prometheus acts purely as an HTTP client: it POSTs compressed Protobuf write requests to the receiver. The remote system is a sink — it never influences scraping, rule evaluation, or alerting on the local server.
A representative configuration in prometheus.yml (Prometheus 2.x, run as the prometheus user with permission to reload the config):
remote_write:\n - url: \"https://remote-write.example.internal/api/v1/write\"\n basic_auth:\n username: prometheus\n password_file: /etc/prometheus/secrets/rw_password\n tls_config:\n ca_file: /etc/prometheus/secrets/ca.crt\n queue_config:\n capacity: 10000\n max_shards: 50\n max_samples_per_send: 5000\n write_relabel_configs:\n - source_labels: [__name__]\n regex: \"go_.*\"\n action: dropTwo decisions are embedded here. First, write_relabel_configs drops low-value series (Go runtime metrics in this example) before they leave the box — filtering at the source is the cheapest way to control remote storage cost and cardinality. Second, queue_config sizes the in-memory shard queue; the defaults are conservative and usually need tuning only after you observe backlog metrics.
Apply with a config reload (curl -X POST http://localhost:9090/-/reload requires --web.enable-lifecycle; otherwise send SIGHUP or restart). Risk: a malformed remote_write block prevents startup on restart but a reload simply fails and keeps the old config — check prometheus_config_last_reload_successful after every change.
Trust and data boundaries
The trust boundary sits at the remote endpoint. Use TLS with a CA you control (or a public CA for managed services) and authenticate with basic auth, bearer tokens, or mTLS depending on what the receiver supports. Credentials in password_file should be readable only by the prometheus user.
Data-wise, assume everything scraped crosses the boundary unless explicitly relabeled away. If some metrics are sensitive (labels containing user IDs, internal hostnames), drop or rewrite them in write_relabel_configs, not at the receiver — by then they have already left your network. Also note the consistency boundary: Remote Write is best-effort. Samples can be duplicated at the receiver after retries, and there is no end-to-end acknowledgment protocol guaranteeing delivery. Deduplication and gap handling are the remote store's problem, which is why backends like Thanos and Cortex implement dedup logic.
Operational checks
Prometheus exposes its own remote-write health as metrics. Alert on and dashboard at least these:
rate(prometheus_remote_storage_samples_total[5m])— confirms data is flowing; a drop to zero while scraping is healthy means the pipeline is broken.rate(prometheus_remote_storage_samples_failed_total[5m])andprometheus_remote_storage_batches_failed_total— non-zero rates indicate the receiver is rejecting or unreachable.prometheus_remote_storage_queue_highest_sent_timestamp_secondscompared against current time — a growing lag means the queue is falling behind.prometheus_remote_storage_shards_desired / prometheus_remote_storage_shards— if desired exceeds actual for sustained periods, raisemax_shards.
End-to-end verification: query the remote store for a known metric (for example up) and confirm recent timestamps. Do not assume config success equals data arrival.
Failure modes
When the remote store is down, Prometheus keeps unsent samples in the WAL on disk and retries with backoff. Short outages (minutes to a few hours, depending on disk headroom and ingestion rate) heal automatically. The dangerous cases:
- Extended outage: the WAL grows until it hits the retention or disk limit, after which the oldest buffered data is truncated — silent remote-side gaps.
- Memory pressure: shard queues hold samples in memory; an undersized
capacitycauses backpressure, an oversized one inflates RSS during catch-up after an outage. - Cardinality explosion: a bad label (user ID, request ID) multiplies series and can overwhelm the remote backend regardless of how well Prometheus itself is tuned. No queue setting fixes this; only relabeling or fixing the instrumentation does.
- Permanent rejection: a receiver returning 4xx (for example on out-of-order or invalid samples) causes those batches to be dropped rather than retried — watch
samples_failed_totalseparately from 5xx-style retries.
A practical failure drill: block egress to the remote endpoint with a firewall rule for 30 minutes in a staging environment, observe WAL directory growth (du -sh /var/lib/prometheus/wal) and the queue lag metric, then restore connectivity and confirm catch-up completes without OOM.
When to change the design
Revisit the architecture when any of these become true:
- Query volume shifts to historical data. If users routinely query beyond local retention, add Remote Read or (more commonly) point Grafana at the remote store's query endpoint. Remote Read pulls data back through Prometheus and adds latency and load to the local server; querying the store directly scales better.
- One Prometheus is no longer enough. Multiple Prometheus replicas writing to the same store require external labels (
external_labelswith a unique replica identifier) so the backend can deduplicate. - Delivery guarantees matter. If gaps are unacceptable, Remote Write alone is the wrong tool — you need replicated Prometheus pairs with backend dedup, or a different ingestion path entirely.
The stable core of the design is simple: Prometheus scrapes and alerts locally, exports best-effort to a sink, and you monitor the export pipeline as a first-class system. Everything else — Remote Read, HA pairs, dedup — is a response to a measured requirement, not a default.
0 replies
A thoughtful contribution can make all the difference. Be the first to share one.