Most OpenTelemetry deployments start as a working demo: SDKs export OTLP to a collector, the collector exports to a backend, and traces appear. The pipeline then meets production traffic, a backend outage, a deploy that doubles span volume, or a noisy tenant, and data starts disappearing without any error anyone sees. The reason is structural. A telemetry pipeline is a chain of buffers, and every buffer has a capacity, a policy for what happens when it fills, and a retry behaviour. Pipeline design means choosing those deliberately.
This article treats the pipeline as an engineering system rather than a configuration file. It assumes you know what the collector is (the collector architecture article covers that) and focuses on processor ordering, backpressure, durable queues, sizing with a worked example, loss accounting, isolation between signals and tenants, and safe operation.
Where data waits and where it dies
Follow one span from creation to storage. In the application, the SDK's batch span processor puts finished spans into an in-memory queue and exports them in batches. The OpenTelemetry specification's defaults for that processor are a queue of 2,048 spans, a 5-second schedule delay, export batches of up to 512 spans and a 30-second export timeout, configurable through variables such as OTEL_BSP_MAX_QUEUE_SIZE. When the queue is full, new spans are dropped. The application keeps running, which is correct, but the loss is visible only in SDK metrics or logs, if at all.
Next, an agent collector on the same node receives the OTLP request. If it is overloaded it can refuse data, and whether that turns into a retry or a loss depends on the sender. The agent forwards to a gateway tier, which forwards to the backend; each exporter has its own queue and retry policy. The backend itself may throttle with HTTP 429 or gRPC RESOURCE_EXHAUSTED. So there are at least four places a span can wait and four places it can be dropped, and a design should state, for each, the capacity and the behaviour when that capacity is exceeded.
Topology: choose the fewest hops that give you the controls you need
Exporting directly from SDKs to a backend is the simplest topology and acceptable for small systems, but it places credentials, endpoints and retry behaviour in every application. An agent per node (a Kubernetes DaemonSet or a sidecar) gives applications a local, cheap endpoint, adds host and Kubernetes metadata, and absorbs short spikes. A gateway tier (a horizontally scaled deployment behind a load balancer) centralises credentials, redaction, routing and cost controls. Tail sampling adds one more requirement: all spans of a trace must reach the same instance, so the agents or a first gateway layer use the loadbalancing exporter with routing_key: traceID to send each trace to a consistent sampling instance. The tail sampling article covers that layer.
Each hop adds latency, cost and another buffer to size. Add one only when it buys a control you need: the agent buys local enrichment and spike absorption, the gateway buys central policy, and the load-balancing layer buys trace-complete sampling.
Processor order inside a pipeline
Within a collector, a pipeline is receivers, then processors in the order listed, then exporters. Order matters because each processor costs memory and CPU, and work done on data that is later dropped is wasted. A sound default order is:
- memory_limiter first, so the collector refuses new data before it runs out of memory, rather than after processing it.
- Filtering and sampling early: drop health-check spans, debug logs and unused metrics before spending effort on them.
- Redaction and attribute hygiene before anything leaves the trust boundary, and before data reaches any exporter.
- Enrichment (for example Kubernetes metadata) after filtering, since it adds bytes to every item it touches.
- Batching last, so exporters send large, compressed requests.
extensions:
file_storage:
directory: /var/lib/otelcol/queue # persistent volume, not emptyDir
processors:
memory_limiter:
check_interval: 1s
limit_percentage: 80 # hard limit, relative to container memory
spike_limit_percentage: 20 # soft limit = 80% - 20% = 60%
filter/health:
error_mode: ignore
traces:
span:
- 'attributes["url.path"] == "/healthz"'
attributes/redact:
actions:
- key: http.request.header.authorization
action: delete
batch:
send_batch_size: 8192 # spans per batch, so one queue entry is about 8192 spans
exporters:
otlp/backend:
endpoint: traces.example.internal:4317
timeout: 10s
retry_on_failure:
enabled: true
initial_interval: 5s
max_interval: 30s
max_elapsed_time: 600s # keep retrying through a 10-minute outage
sending_queue:
enabled: true
num_consumers: 10
queue_size: 1000 # batches (default requests sizer); see sizing below
storage: file_storage # survive restarts; bounded by disk
service:
extensions: [file_storage]
pipelines:
traces:
receivers: [otlp]
processors: [memory_limiter, filter/health, attributes/redact, batch]
exporters: [otlp/backend]Two notes on this configuration. The memory limiter's soft limit is the hard limit minus the spike allowance; above the soft limit it refuses new data, and above the hard limit it also forces garbage collection. Its check interval matters, because memory can grow between checks; one second is the documented recommendation. And recent collector versions can batch inside the exporter's sending queue (a batch block under sending_queue), which you may prefer to the batch processor because batching then happens after the durable queue rather than before it. Check your collector version's documentation and pick one approach.
Backpressure and durability
Backpressure means telling the sender to slow down instead of silently dropping. When the memory limiter refuses data, an OTLP receiver returns an error and the upstream sender can retry. That is only useful if the sender does retry and has somewhere to hold data meanwhile; the memory limiter documentation warns that data is permanently lost if the component before it does not retry. In a two-tier design, the agent's exporter queue is that somewhere. At the edge, the SDK's queue is small by design, so backpressure ending at the application means drops.
The exporter's sending_queue is the main buffer. Its defaults are worth knowing because they often surprise people: it is enabled, holds 1,000 entries, uses 10 consumers, and its sizer defaults to requests, so capacity depends on how big each batch is. Setting sizer: items or bytes makes an in-memory limit predictable; the documentation describes the persistent queue's limit in batches, so with storage enabled, size it in batches and keep batch sizes stable. When the queue is full, new data is rejected unless block_on_overflow is set, which pushes back upstream instead. Retries use exponential backoff starting at 5 seconds, capped at 30 seconds between attempts, and give up after 300 seconds in total by default, after which the batch is dropped.
The in-memory queue is lost if the collector restarts, which during an incident is precisely when restarts happen. Setting storage to a storage extension such as file_storage makes the queue persistent, so it survives restarts and can be bounded by disk rather than memory. The trade-offs are disk I/O on every item, the need for a persistent volume per instance, and a drain period after recovery during which fresh data competes with the backlog.
Worked example: sizing a gateway tier
Suppose a cluster produces 40,000 spans per second at peak, averaging 1 KB each as serialised OTLP before compression. The requirement is to survive a 10-minute backend outage without losing data, and to recover within 10 minutes afterwards.
- Throughput: 40,000 spans/s x 1 KB = about 40 MB/s, or 2.4 GB per minute, entering the gateway tier.
- Outage backlog: 10 minutes x 2.4 GB = about 24 GB, or 24 million spans. That is far more than you should hold in collector memory, so the design needs persistent queues: with six gateway replicas, each needs roughly 4 GB of queue storage plus headroom, so provision about 8 GB per volume. At 8,192-span batches, 4 million spans is about 500 batches, so
queue_size: 1000gives headroom for partial batches flushed on timeout. - Retry budget:
max_elapsed_timemust exceed the outage you plan to survive, so raise it from the default 300 seconds to at least 600. - Recovery: to drain 24 GB in 10 minutes while new data keeps arriving at 40 MB/s, the tier must export at about 80 MB/s, twice peak. Confirm the backend will accept that rate rather than throttling you, and size gateway CPU for it.
- Memory: set the container memory limit from a load test at twice peak with the memory limiter enabled, not from a guess; keep the limiter's hard limit below the container limit so the collector refuses data before the kernel kills it.
If the storage or the 2x drain rate is too expensive, decide explicitly what to drop. A common answer is to keep error and slow traces durable while sampling the rest harder during backlog, which is a policy decision you should write down rather than leave to whichever queue fills first.
Loss accounting with self-telemetry
A pipeline you cannot account for will lose data quietly. The collector reports its own metrics; scrape them like any other service. The ones to dashboard and alert on are, per component: items accepted and refused by receivers (otelcol_receiver_accepted_spans, otelcol_receiver_refused_spans), items sent and failed by exporters (otelcol_exporter_sent_spans, otelcol_exporter_send_failed_spans), items that could not be queued (otelcol_exporter_enqueue_failed_spans), and queue depth against capacity (otelcol_exporter_queue_size and otelcol_exporter_queue_capacity). Equivalent metrics exist for metric points and log records. Exact names can gain suffixes such as _total depending on the collector version and how metrics are exposed, so confirm them against your own scrape.
Then reconcile end to end. Compare spans produced (from SDK metrics, or a span-count metric derived in the agent) with spans stored by the backend, per service, and alert when the gap exceeds a threshold such as 1% for ten minutes. Alert on queue utilisation above 50% sustained, any rise in enqueue failures, and refused data at receivers. Treat a sustained loss the way you treat an SLO burn: it means your incident data is incomplete precisely when you need it. Metrics derived from spans, as in the span metrics article, should be computed before sampling so loss and sampling do not distort RED dashboards.
Isolation between signals and tenants
Run separate pipelines per signal so a log flood cannot starve traces: each pipeline has its own exporters and queues, even when they point at the same backend. For multi-tenant platforms, route by a resource attribute or request metadata into per-tenant or per-tier pipelines, give each its own queue and limits, and apply rate limits or sampling to the noisiest. Cardinality is the other tenant risk: one team adding a user ID as a metric attribute can multiply series and cost. Enforce attribute allowlists in the gateway, as described in the cardinality article.
Failure modes
| Failure | What you see | Design response |
|---|---|---|
| Backend outage longer than retry budget | Gap in traces after recovery | max_elapsed_time above planned outage; persistent queue |
| Collector OOM-killed | Restarts, lost in-memory queue | memory_limiter first; hard limit below container limit |
| SDK queue overflow during traffic spike | Missing spans from busy services only | Local agent endpoint; tune OTEL_BSP_ settings; watch SDK drop metrics |
| Tail sampler receives partial traces | Wrong keep decisions, broken traces | Load-balance by trace ID; stable membership during deploys |
| One signal or tenant floods the pipeline | Other signals delayed or dropped | Separate pipelines and queues; per-tenant limits |
| Bad config deployed fleet-wide | All telemetry stops at once | Validate config in CI; canary collector rollouts |
Operating the pipeline
Treat collector configuration as code: review it, validate it in CI with the collector's validate command against the exact binary you deploy, and roll it out to a canary slice first while watching the self-telemetry above. Build a custom collector with the OpenTelemetry Collector Builder containing only the components you use, which shrinks the attack surface and makes upgrades deliberate. Pin versions and read release notes, since component configuration does change between releases. Rehearse a backend outage in staging at least once: block the exporter endpoint for longer than your planned outage, confirm the queues fill and drain as calculated, and fix whatever the rehearsal reveals.
What to do next
- Draw your pipeline and list every buffer with its capacity, overflow behaviour and retry policy, from SDK queue to backend.
- Put memory_limiter first in every pipeline and set its hard limit below the container memory limit.
- Size sending_queue from a written outage target (in batches when it is persistent), and enable file_storage on a persistent volume where data must survive restarts.
- Raise retry_on_failure max_elapsed_time above the outage you plan to survive, and confirm the backend accepts a 2x drain rate.
- Dashboard accepted, refused, sent, failed and enqueue-failed counts plus queue utilisation, and alert on end-to-end produced-versus-stored gaps.
- Separate pipelines per signal, validate config in CI, canary every rollout, and rehearse a backend outage in staging.