Skip to content
All writing

At-least-once delivery without losing sleep

, 5 min read

Metrics pipelines rarely fail loudly. The usual failure is quieter: the database slows down, a consumer times out, and a few minutes of data never arrive. Nobody notices until a graph has a gap in it during an incident.

When I built my Datadog-style observability pipeline, I wanted every one of those failure modes to be a decision rather than an accident.

The shape of the pipeline

A Go service (Gin) receives metrics over HTTP, validates them and publishes to a Kafka topic with 6 partitions. Anything malformed goes to a dead-letter queue instead of being dropped. A Rust consumer built on tokio reads the topic, runs EWMA and rolling Z-score anomaly detection, and writes to ClickHouse.

Commit after the write, not before

Kafka consumers track their position with offsets. The tempting default is to auto-commit them on a timer. That's the bug: if the consumer commits offset 1,000 and then the ClickHouse write for messages 990 to 1,000 fails, those messages are gone. Kafka thinks you have them.

So the consumer commits an offset only after the batch containing it has been written successfully. Simplified, the loop looks like this:

let batch = collect_batch(&mut stream, MAX_BATCH, MAX_WAIT).await;
clickhouse.insert(&batch).await?;   // if this fails, we return early
consumer.commit(&batch.last_offset()).await?;  // only now is it "done"

If the process crashes between the insert and the commit, the batch is replayed on restart. That's what at-least-once means: duplicates are possible, loss is not. Duplicates are a much easier problem, because ClickHouse rollups and idempotent keys can absorb them.

Don't hammer a struggling database

Committing after writes creates a new risk. If ClickHouse is degraded, every batch fails, the consumer retries immediately, and you have now built a very efficient way to make a slow database slower.

The consumer loop is wrapped in a circuit breaker:

  • After 5 consecutive failures the breaker opens and the consumer stops writing.
  • It retries with exponential backoff, capped at 60 seconds.
  • One successful write closes it again.

While the breaker is open, messages wait safely in Kafka. The visible symptom is consumer lag, which is bounded, measurable and alertable. That's the whole trick: turn "we lost data" into "we were late", and make late show up on a dashboard.

Proving it

Claims like this are cheap, so the project ships with 53 unit tests (41 in Go, 12 in Rust), 25 integration tests and a k6 load-test script for the full 14-container stack. The integration tests are the important ones, because delivery guarantees only mean something when the whole stack is running.

If you're building something similar, the one-line takeaway: decide where your data is allowed to wait, and make sure that's the only place it can go when things break.