Appears in: Telemetry Ingestion Pipeline §1 (consistency and durability requirements), §3.1 (agent retry/backoff, gateway ACK), §3.2 (Kafka producer/consumer semantics).
This looks like a question about delivery semantics — at-least-once vs exactly-once. It’s actually a retry-policy question wearing a delivery-semantics costume. Every hop in this pipeline answers the same three questions — is this failure worth retrying, how long before the next attempt, and when do I give up — before “at-least-once” or “exactly-once” ever enters the picture. The delivery semantic a pipeline ends up with is the accumulated side effect of those retry decisions across every hop, not an independent design choice you make separately.
The uncertain moment
Producer sends message ──▶ [ network / broker ] ──▶ ??? did it arrive?
│
Producer gets a timeout, not an answer.
Did the message get lost? Or did it arrive and
only the ACK got lost on the way back?
The producer cannot tell the difference.
The producer has exactly two choices when it can’t confirm delivery: retry (risking a duplicate, if the original actually did arrive) or don’t retry (risking silent loss, if it didn’t). Every retry policy is a set of rules for which side of that coin you accept, and under what conditions.
Anatomy of a retry policy — three knobs
A retry policy is not one decision, it’s three, and getting any one of them wrong produces a specific, recognizable failure mode:
| Knob | Question it answers | What gets it wrong |
|---|---|---|
| Trigger | Which failures are worth retrying? | Retrying a 400 Bad Request forever (it will never succeed); not retrying a 429/503 that would succeed next time |
| Backoff | How long before the next attempt? | Fixed-interval retries from thousands of clients synchronize into a thundering herd on the exact system that’s struggling |
| Budget | When do you stop? | No cap → a permanently-failing message or dependency retries forever, blocking everything queued behind it |
flowchart TD
A["Call fails"] --> B{"Retryable?\n(timeout, 429, 503 — yes\n400, 422 — no)"}
B -->|"No — permanent failure"| FAIL["Fail fast\n(poison-pill / DLQ path)"]
B -->|"Yes — transient"| C{"Within retry budget?\n(max attempts / max elapsed time)"}
C -->|"No — budget exhausted"| GIVEUP["Give up\nDLQ, orphan handling, or drop"]
C -->|"Yes"| D["Backoff\n(exponential + jitter)"]
D --> E["Retry"]
E -->|"Success"| DONE["Done — but did the\noriginal attempt also land?\n→ duplicate risk"]
E -->|"Fails again"| B
That last edge — “did the original attempt also land?” — is where delivery semantics come from. It’s not a separate branch in the flowchart; it’s the question the whole retry loop leaves unanswered every time it takes the “retry” path.
Backoff strategy: why exponential + jitter, not fixed-interval
Fixed-interval retry (always wait exactly N seconds) is the naive default, and it fails in a specific way: when the thing that made everything fail simultaneously — a gateway rolling restart, a regional outage — resolves, every client that was waiting retries at the same instant. That’s a thundering herd, and it’s the gateway fan-in problem at 100K+ agents: 50K agents reconnecting simultaneously after a rolling restart.
Exponential backoff (1s → 2s → 4s → … → capped at 60s, per the agent WAL config in Q1) spreads out load over time instead of hammering the dependency at a constant rate while it’s still recovering.
Jitter (randomizing each wait within a range, rather than a deterministic sequence) is what actually breaks the synchronization — without it, exponential backoff alone still leaves every client on the same schedule, just a slower one. The same fix shows up in an unrelated subsystem for the same underlying reason: Q6 jitters ingester flush schedules across tenants specifically because ingesters started at the same time otherwise resynchronize their 2-hour flush boundaries over time — identical failure shape (correlated retries becoming a self-inflicted storm), different subsystem (scheduled flush, not a failed RPC).
# Alloy agent config — exponential backoff + jitter + a bounded retry budget
retry_on_failure:
enabled: true
initial_interval: 1s
max_interval: 60s # cap — without this, backoff grows unbounded
max_elapsed_time: 120s # budget — after this, give up and let the WAL hold the batch
retry_on_http_429: true # treat backpressure as retryable, not a permanent failure
Retry budgets and giving up
A retry policy needs a stopping condition on two different axes:
- Max attempts / max elapsed time —
delivery.timeout.ms = 120sin the Kafka producer config (Q1) bounds how long a single message can spend retrying before the producer gives up on it entirely. - What happens at that boundary depends on what’s failing:
- One bad message (malformed payload that will never deserialize) is a poison pill (Q4) — retrying it is pointless because the failure is permanent, not transient, and an unbounded retry blocks every message queued behind it on that partition. The fix is a retry-count cap per message, with overflow routed to a dead-letter destination instead of retried forever.
- One dependency being down (not one message) calls for a circuit breaker: stop sending traffic to a dependency that’s already failing, so retries stop amplifying the outage and the dependency gets room to recover. This is the same “stop trying, protect the shared resource” pattern Q2 uses at the processor layer — there it’s a cardinality-budget reject rather than a downstream-failure trip, but it’s the same shape: fail fast at a cheap layer rather than let the failure propagate and get expensive downstream.
A retry policy that conflates these two — retrying a poison pill as if it were a transient failure, or retrying against a dead dependency as if one more attempt might get lucky — is the single most common root cause of “why is consumer lag growing on exactly one partition” or “why did a 10-minute outage turn into an hour of recovery.”
Retry storms and the backpressure feedback loop
Retries aren’t free even when they eventually succeed — every retry is additional load on a system that, by definition, just failed to keep up. At scale this becomes a feedback loop the pipeline has to close deliberately, not just tolerate:
flowchart TD
A["Storage full"] --> B["Processor slows\nconsumer group lag grows"]
B --> C["Kafka consumer lag alarm fires\n→ scale out processors via HPA"]
C --> D{"Lag persists?"}
D -->|No| DONE["Normal operation resumes"]
D -->|Yes| E["Gateway returns gRPC RESOURCE_EXHAUSTED"]
E --> F["Agent receives 429\n→ backs off with jitter, doesn't retry harder"]
F --> G["Agent local buffer absorbs burst\nWAL / memory queue"]
The critical design choice: 429/RESOURCE_EXHAUSTED is a signal to slow the retry rate down,
not a failure to retry around faster. A retry policy that treats backpressure as “just another
transient failure, retry immediately” turns a load-shedding signal into more load — exactly the
scenario exponential backoff with jitter exists to prevent. This is also why the gateway itself must
never become the coordination point for backpressure
(main design) — the decision to slow down has to live in the
retrying client, where the backoff state actually is.
The side effect of retrying: delivery semantics
A retry policy whose trigger rule is “always retry on uncertainty” produces at-least-once delivery by construction — not because anyone chose at-least-once as a semantic, but because duplicates are the unavoidable consequence of that trigger rule. Delivery semantics are downstream of the retry decision, not upstream of it:
| Semantic | Retry policy | What can go wrong | Cost |
|---|---|---|---|
| At-most-once | Never retry | Message silently lost | Cheapest — fire and forget |
| At-least-once | Always retry on uncertainty | Message delivered twice (or more) | Cheap — consumer must handle duplicates |
| Exactly-once | Retry, but the system guarantees no duplicate and no loss | Nothing, by definition — but achieving this is the expensive part | Highest — coordination overhead on every message |
flowchart TD
A["Producer sends, times out waiting for ACK"] --> B{"Retry policy trigger"}
B -->|"Never retry"| AM["At-most-once\nmessage may be lost"]
B -->|"Always retry"| AL["At-least-once\nmessage may be duplicated"]
AL --> DEDUP{"Consumer dedups\nby a stable ID?"}
DEDUP -->|Yes| EFF["Effectively-once\n(the practical version of exactly-once)"]
DEDUP -->|No| DUP["Duplicate side effects\n(double-counted metric, duplicate charge, etc.)"]
At-most-once is rarely chosen deliberately for anything that matters — it’s what you get by default if you don’t build retry logic at all. It shows up in fire-and-forget protocols like StatsD over UDP, where the cost of guaranteeing delivery exceeds the value of any single sample.
At-least-once is the default for almost everything else, because retrying on uncertainty is the safe default when duplicates are cheaper to handle than loss. The burden shifts downstream: whatever consumes the message now has to tolerate seeing it more than once.
Why “exactly-once” is expensive — the actual mechanism
There is no way for a producer to make a single network call that is guaranteed to have exactly one effect on a remote system it can’t see the internal state of — this is fundamentally the same shape of problem as the Two Generals’ Problem in distributed systems theory. What “exactly-once” systems actually do is fake the outcome through one of two mechanisms — both of which exist specifically to absorb the duplicates a retry policy creates:
Mechanism 1: Idempotent producer (dedup at the write). The producer attaches a unique, monotonically increasing sequence number (per producer session) to every message. The receiver keeps track of the last sequence number it accepted per producer and silently discards anything it’s already seen — a retry of an already-delivered message becomes a no-op instead of a duplicate.
# Kafka producer config — this is exactly-once at the produce step only
enable.idempotence: true # broker deduplicates retries by (producer_id, sequence_number)
acks: all # all ISR replicas must confirm before considering it sent
This solves duplication on the hop between producer and broker — it does not, by itself, make the entire end-to-end pipeline exactly-once. A message can still be processed twice further downstream if the consumer crashes after processing but before committing its offset — a retry at a different hop, with the same root cause.
Mechanism 2: Transactional consume-process-produce (dedup across a boundary). Kafka’s transactional API lets a consumer atomically commit “I processed message X” (its consumer offset) together with “I produced message Y” (its output) as a single all-or-nothing unit. If the consumer crashes between processing and committing, the whole transaction aborts and gets replayed from the start — so from an outside observer’s point of view, either both things happened or neither did.
Transaction boundary:
┌─────────────────────────────────────┐
│ 1. Consume message X (offset N) │
│ 2. Process X → produce message Y │
│ 3. Commit offset N + produce Y │ ← atomic: all or nothing
└─────────────────────────────────────┘
The catch that matters most in an interview: this only guarantees exactly-once within Kafka’s own transactional boundary. The moment the pipeline writes to something outside that boundary — an external database, an external API, a metric TSDB — the guarantee doesn’t automatically extend there. The external system needs its own idempotency (e.g., an upsert keyed by a unique ID) or the guarantee silently stops at the Kafka boundary, and any retry past that point is back to producing plain at-least-once duplicates.
”Effectively-once” — what most systems actually build
Because true end-to-end exactly-once requires every hop in the chain (including external systems) to participate in the same transactional or idempotency scheme, most real systems settle for:
Effectively-once = at-least-once retry policy + idempotent processing at the consumer
The message can arrive more than once; the consumer recognizes duplicates (by a stable ID, content hash, or fingerprint) and processing a duplicate has no additional effect. This achieves the outcome users care about (no double-counting, no duplicate charges) without needing the expensive coordination machinery of a literal exactly-once guarantee.
Retry policy per hop, in this pipeline
Retry policy isn’t set once for the whole pipeline — it’s set per hop, and the duplicate risk compounds across every hop that retries independently:
| Hop | Retry trigger | Backoff | Budget | Duplicate risk | Mitigation |
|---|---|---|---|---|---|
| Agent → Gateway | No ACK / timeout | Exponential + jitter (1s → 60s) | WAL max_age (bounded local buffer) | Yes — ambiguous ACK | Idempotency key + gateway-side dedup, tenant-scoped for the billing tenant |
| Gateway → Kafka | Producer timeout | Kafka producer internal backoff | delivery.timeout.ms = 120s | Yes | enable.idempotence: true — broker dedups by (producer_id, sequence_number) |
| Kafka consume → Processor | Consumer crash before offset commit | Consumer group rejoin, resumes from last committed offset | Consumer group session timeout | Yes — reprocesses the batch | Accepted at-least-once; downstream store dedup (below). Billing tenant uses transactional consume-process-produce |
| Processor → Storage (Mimir/Loki/Tempo) | Remote-write timeout | Exponential backoff | Bounded retry queue | Yes | Idempotency key + storage-side dedup, tenant-scoped for the billing tenant |
Full trace of how this plays out for the one tenant where duplicates aren’t tolerable: Q12.
Which signals can tolerate the duplicates at-least-once produces
| Signal | Retry policy | Why it’s safe | | ----------------------------------------- | ----------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------- | --------------------------------------------------------------------------------------------------- | | Metrics | At-least-once | The TSDB deduplicates by timestamp + label fingerprint — a duplicate sample at the same timestamp is simply overwritten, not double-counted | | Logs | At-least-once | Dedup on a content hash + timestamp window within the log store | | Traces | At-least-once | Trace stores deduplicate by span ID — a duplicate span is a no-op | | Billing-critical metrics (exception case) | Exactly-once (scoped to that tenant only) | See Q12 — the cost is only justified because a dollar amount depends on the count being exactly right |
The main design’s stance, stated directly: default every hop to always-retry (at-least-once), and only pay for exactly-once where there is a billing or compliance requirement. Everywhere else, at-least-once plus a cheap dedup key at the storage layer gets you the same practical correctness at a fraction of the latency and operational cost — paying for transactional coordination on every metric sample, for a guarantee the TSDB’s own dedup already provides for free, would be solving a problem that doesn’t exist.
What retry policy choices cost
This is the number to have ready when someone asks “why not just always retry aggressively” or “why not just always do exactly-once”:
- Latency — transactional commits require additional coordination round trips per message (or per batch) compared to fire-and-forget or simple at-least-once acknowledgment.
- Throughput — transaction coordinators serialize commits; you trade raw throughput for the correctness guarantee.
- Operational surface area — a dedicated consumer group and transactional path, isolated from the shared processing fleet’s default at-least-once path, per Q12‘s answer — because you cannot mix transactional and non-transactional consumption in the same consumer group.
- Blast-radius amplification — a retry policy that doesn’t back off (or doesn’t jitter) turns a transient dependency slowdown into a self-inflicted overload. Q7‘s 10-minute regional outage stays a 10-minute outage specifically because agents back off with jitter and buffer locally instead of retrying harder against an endpoint that’s already down.
- False confidence risk — a bug in the exactly-once mechanism itself still produces internally-consistent-looking metrics. The only real check is reconciliation against a source of truth outside the pipeline, not a pipeline-internal counter.
Related
- Telemetry Ingestion Pipeline (full design) — §1 (consistency/durability requirements), §3.1 (agent retry/backoff, gateway ACK), §3.2 (Kafka producer config and exactly-once vs at-least-once)
- Q1: 500M ingest, zero-drop rolling deploy —
source of the WAL + exponential backoff config and
delivery.timeout.msbudget - Q2: Cardinality storm detection & mitigation — the circuit-breaker pattern that pairs with retry budgets to protect a shared resource
- Q4: Metric point journey failure points — poison-pill messages: the retry-budget failure mode at the single-message level
- Q6: Compactor storm diagnosis — the same correlated-retry-becomes-a-storm failure shape, in scheduled flushes rather than failed RPCs
- Q7: Regional gateway outage blast radius — why backoff + jitter is what keeps an outage’s blast radius from growing
- Q12: Mixed exactly-once billing tenant — the one scenario in this design where exactly-once is actually worth the cost
- 05 — Backpressure — the retry behavior that produces at-least-once duplicates is the same mechanism backpressure-driven agent retries rely on
Local graph
Linked from 7 notes
3.2 Layer 2: Durable Buffer (Kafka)
Layer 2 of the telemetry ingestion pipeline: the Kafka durable buffer — topic design, partitioning strategy, hot-spots, retention, retry/delivery semantics, producer config, consumer lag, and schema evolution.
Q1: 500M Samples/Sec, Zero Drop on Rolling Deploy
Full principal-level solution: design a telemetry ingestion pipeline for 500M metric samples/sec from 100K services globally with a zero-drop guarantee during rolling deployment of the ingestion tier.
Q12: Exactly-Once for One Billing Tenant While Others Stay At-Least-Once
Full principal-level solution: support exactly-once ingestion for a single billing-critical tenant in a shared pipeline that is at-least-once everywhere else, and account for what it costs.
9 — OTel Collector Pipeline Design
Receivers, processors, and exporters chained into a pipeline; why a platform runs more than one; and the agent/gateway topology that tail sampling specifically forces on that design.
1. Clarify Requirements First
The first-5-minutes clarifying questions for the telemetry ingestion pipeline design — signal types, scale envelope, consistency/durability, multi-tenancy, and protocol — whose answers change the entire architecture.
3.1 Layer 1: Ingestion Frontier
Layer 1 of the telemetry ingestion pipeline: the ingestion frontier — responsibilities, fan-in at 100K+ agents, protocol negotiation, batching, backpressure, and rate limiting.
System Design
Principal/Staff-level system design reference collection for MAANG interview preparation — observability pipelines, distributed systems, reliability engineering, and beyond.
Related notes
3.2 Layer 2: Durable Buffer (Kafka)
Layer 2 of the telemetry ingestion pipeline: the Kafka durable buffer — topic design, partitioning strategy, hot-spots, retention, retry/delivery semantics, producer config, consumer lag, and schema evolution.
Chapter 1 — Telemetry Ingestion Pipeline
Principal/Staff-level design of a high-throughput telemetry ingestion pipeline — requirements, architecture, deep dives, and trade-offs at 10x scale.
Q1: 500M Samples/Sec, Zero Drop on Rolling Deploy
Full principal-level solution: design a telemetry ingestion pipeline for 500M metric samples/sec from 100K services globally with a zero-drop guarantee during rolling deployment of the ingestion tier.
Q10: Self-Service Tenant Onboarding With Zero Platform-Team Involvement
Full principal-level solution: design a self-service tenant onboarding API for a telemetry pipeline that protects shared infrastructure from a misbehaving new tenant on day one.