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

OpenTelemetry Collector Resilience: Memory, Queues, Retries, and Backpressure

Engineer Collector memory limits, batch processors, exporter queues, retries, and persistent storage with explicit backpressure, loss, and duplication tradeoffs.

Collector resilience is a policy for what happens when telemetry arrives faster than it can be processed or a downstream backend stops accepting data. A memory limiter, a batch processor, an exporter sending queue, retry policy, and persistent storage each protect a different boundary. Enabling all of them does not create an exactly-once delivery system; it creates a bounded pipeline whose loss, latency, memory, disk, and backpressure behavior must be measured.

The right objective is not “never drop telemetry.” It is to state which signals may be delayed or lost, for how long, what the Collector does when each bound is reached, and which upstream senders can recover. A metrics scrape source may not replay a missed interval, while an OTLP client may retry a rejected request. Those are different failure contracts even when both feed the same Collector.

Understand the stages before choosing limits

Incoming receivers first allocate and decode data, then pass it through ordered processors. The memory limiter should be early enough to return pressure to a receiver before downstream processors accumulate more state. A batch processor holds data to form larger sends. Exporter queueing then buffers work for one exporter, and exporter retry handles eligible send failures for batches already being attempted. Persistent queue storage can preserve queued work across a Collector process restart when the exporter and storage extension support it.

These are not interchangeable controls. Batching is an efficiency and latency decision; it is not a durable buffer. A queue absorbs a finite mismatch between incoming and outgoing rates; it is not infinite retry. Retry cannot recover data that was rejected before reaching the queue. The Collector’s current resiliency guidance explicitly calls out queue overflow, retry expiry, Collector restart, storage failure, and configuration errors as loss paths.

Put memory protection at the front of each pipeline

The memory_limiter periodically measures process-heap use. It has a hard limit and a soft threshold calculated as limit_mib - spike_limit_mib. Above the soft threshold it refuses incoming data with a non-permanent error and forces garbage collection at the hard limit. Correctly implemented upstream receivers should retry the same data and may apply backpressure; a receiver or source that cannot retry can lose data. An exporter queue retaining live objects during a backend outage can also prevent garbage collection from reclaiming memory.

This processor is a guardrail, not a hard cgroup enforcement mechanism or a substitute for sizing. Its checks are periodic, incoming data can consume memory before it reaches the processor, and total process use can exceed heap allocation. Leave headroom below the container’s memory limit for runtime and non-heap allocations. The upstream project recommends placing the processor first and coordinating it with Go’s GOMEMLIMIT; choose percentages from measured load and the actual container cgroup rather than copying a sample unchanged.

Batch deliberately, then bound the exporter queue

The batch processor sends when send_batch_size is reached or timeout expires. send_batch_size is a trigger, not a maximum. Set send_batch_max_size when the next component needs an enforced upper bound; the processor splits larger batches to meet that limit. Larger batches can reduce per-request overhead, but increase buffering, latency, and the amount of data involved in a failed export. Put batching after drops or sampling processors so work that will be discarded is not needlessly batched.

The example below is a configuration fragment for a Collector distribution that includes otlp receiver, otlphttp exporter, memory_limiter, batch, and the file_storage extension. Replace the endpoint and image/distribution with the versions you actually ship. Configure authentication and TLS policy for your backend; a public endpoint without the required credentials is not a production-ready example. Queue size is deliberately illustrative, not a recommended universal capacity.

receivers:
  otlp:
    protocols:
      grpc:
        endpoint: 0.0.0.0:4317

processors:
  memory_limiter:
    check_interval: 1s
    limit_mib: 768
    spike_limit_mib: 128
  batch:
    timeout: 1s
    send_batch_size: 2048
    send_batch_max_size: 4096

exporters:
  otlphttp/primary:
    endpoint: https://otlp.example.com:4318
    sending_queue:
      enabled: true
      storage: file_storage/queue
      sizer: items
      queue_size: 200000
      num_consumers: 8
      block_on_overflow: false
    retry_on_failure:
      enabled: true
      initial_interval: 5s
      max_interval: 30s
      max_elapsed_time: 5m

extensions:
  file_storage/queue:
    directory: /var/lib/otelcol/queue
    create_directory: true

service:
  extensions: [file_storage/queue]
  pipelines:
    traces:
      receivers: [otlp]
      processors: [memory_limiter, batch]
      exporters: [otlphttp/primary]

When using a fixed 768 MiB heap limit in this example, the soft refusal threshold is 640 MiB. The Collector documentation describes GOMEMLIMIT at 80 percent of the configured hard memory limit as a recommended starting practice; that would be about 614 MiB here, but should still be tested under real load. The container limit must be higher than the processor’s heap target and account for other process memory. If cgroup limits are the sizing source on Linux, the percentage options can decouple the Collector config from the deployment limit, but that does not remove the need to test actual cgroup detection.

Queue units matter. The shared exporter helper can count queue capacity as requests, individual telemetry items, or serialized bytes; exporter implementations may customize whether and how the helper is used. A request count is cheap but one request can be much larger than another. Item counts are easier to reason about but a span, metric point, and log record do not necessarily have the same in-memory cost. Byte sizing is closer to payload size but has extra measurement cost and still does not equal process memory.

If measured throughput is 5,000 records per second and the outage objective is 40 seconds, an item-based queue needs at least 200,000 item slots before adding burst, recovery, and safety margin. That is a capacity calculation, not evidence the Collector can afford that memory. With persistent storage, estimate required disk from the observed serialized bytes per second and tolerated outage, then include compaction headroom and concurrent disk consumers. Validate the estimate during a representative load test.

Queue overflow and retry are separate loss decisions

The exporter queue is bounded. With block_on_overflow: false, queue insertion does not wait for space; if the queue is full or its storage cannot accept more data, an enqueue failure occurs before the exporter retry path. The shared helper reports enqueue-failure metrics, but the exact receiver response and upstream retry behavior depend on the component and protocol. Do not infer that an OTLP client will retry forever because the Collector returned an error.

With block_on_overflow: true, the caller may wait for space until capacity frees or a request times out. This can turn backend pressure into upstream backpressure, which is useful only when upstream producers can tolerate it and propagate it safely. Waiting can tie up receiver work, increase application-side latency, or shift memory pressure to the sender. Test the whole path before enabling it. If a producer has no durable spool or retry contract, a nonblocking rejection and a blocked request have different but still imperfect failure modes.

Retry applies to data the exporter is attempting to send, not every item in the receiver. The common exporter helper documents bounded exponential backoff and a default maximum elapsed retry time; component settings and defaults can differ. After the configured retry budget is exhausted, data may be dropped. Authentication errors, unsupported payloads, and malformed requests are not repaired by repeating them indefinitely. Alert on persistent send errors and fix the destination, credentials, or compatibility problem rather than raising retry time as a substitute for diagnosis.

Ambiguous network outcomes make duplicates possible: the backend may have accepted a batch even though the response was lost, after which a retry sends it again. Retries therefore do not imply exactly-once delivery, and there is no universal Collector-level deduplication contract across signals and exporters. If duplicates are consequential, determine the destination’s idempotency or deduplication behavior and test it under lost-response conditions. For most telemetry, bounded duplicates are preferable to silent gaps, but the choice belongs in the data contract.

Persistent queues reduce one failure window, not all of them

When sending_queue.storage references a configured storage extension, the helper uses persistent queue storage rather than an in-memory queue. file_storage writes queue state to a directory; after a Collector process restarts, the extension can reopen persisted entries and resume export. The directory needs read/write access and should be mounted on storage that survives the restart scenario you intend to protect against. A container’s ephemeral layer, node-local disk, persistent volume, and remote disk have different failure domains.

A local persistent queue is not an independent backup or a broker cluster. Disk exhaustion, filesystem corruption, node loss, a full queue, an exporter retry limit, and an incompatible rollout can still lose or strand data. File storage settings such as fsync, compaction, maximum file size, and permissions affect durability, performance, and failure behavior. The extension’s docs caution that synchronous fsync improves integrity at a performance cost. Treat durable storage as an explicit tested choice, not as a label that proves zero data loss.

Also account for which telemetry transformations happen before the queue. Data removed by filtering or sampling cannot be recovered from the exporter queue. A file queue can contain sensitive telemetry and credentials or resource attributes, so protect the volume, limit access, define retention and cleanup, and include it in disk and backup reviews. Ensure exporters using persisted batches can reacquire current authentication context; the shared helper documentation notes that Auth-extension context is not propagated through persistent queues.

Observe the queue and prove the outage behavior

Monitor queue size and capacity, enqueue failures, send failures, accepted/refused records, process heap/RSS, CPU spent in garbage collection, persistent-storage usage, restarts, and end-to-end telemetry delay. Names and availability of internal metrics vary by Collector version and distribution; check the internal telemetry reference for the exact build. A rising queue is not itself a failure if it drains within the outage objective, but a queue pinned near capacity means the pipeline is losing safety margin.

Run a controlled acceptance test with a known input count and a backend that can be stopped, slowed, then restored. During a short outage, confirm the queue rises without exceeding configured memory and disk budgets; then restore the backend and verify it drains at a measured rate faster than ongoing input. Extend the outage until the planned boundary and verify the documented overflow response, refusal counters, sender retry behavior, and alert. Restart the Collector with queued data and verify recovery from persistent storage. Finally, simulate an ambiguous failed response and inspect the destination for duplicates.

Pin the Collector distribution and version in deployment, because components may be absent or at different stability levels across core, contrib, and vendor builds. Validate the exact config with that binary, load-test the expected burst, and review exporter-specific queue and retry docs before relying on shared-helper defaults. A resilient pipeline is one whose capacity boundaries and failure outcomes are visible and rehearsed, not one with the largest possible queue.

Related:

Sources:

Comments