← All posts
Post 09September 5, 2026

Manual Kafka acknowledgements and a dead-letter path for poison logs

log0's normalization consumer commits offsets only after a successful write or a publish to raw-logs-dlq. The post covers the consume loop, fault injection, and remaining gaps for DLQ replay and alerts.

Ashmit JaiSarita Gupta
Ashmit JaiSarita Gupta

Full-stack Software Engineer - (Builder of log0)

Manual Kafka acknowledgements and a dead-letter path for poison logs

I sent 600 logs through the pipeline. 300 were deliberately poisoned to be unparseable; 300 were clean. Every one of the 600 got a 202. The 300 poison messages landed in raw-logs-dlq, the 300 clean ones flowed through, and the clean path still produced its incident. Zero messages dropped, zero non-202 responses, and not one poisoned record stalled the partition behind it. That outcome rests on a single rule that is easy to state and easy to get wrong: the consumer commits a Kafka offset only after the message is durable somewhere, either processed or quarantined, and never on a try/catch that only logs and moves on.

This is post 9 in a series on building log0. Post 8 ended on the boundary that keeps tenants apart and promised this one: why log0 acknowledges a Kafka message only after it is safely processed or safely quarantined, and what raw-logs-dlq holds when something genuinely cannot be parsed. The short version is that "at-least-once" is a promise you have to keep in code, and the place most pipelines quietly break it is the exception handler.


The bug is the bare catch

Here is the consumer everyone writes first. A message arrives, the consumer processes it, and because processing can throw, it wraps the work:

java
@KafkaListener(topics = "raw-logs")
public void consume(RawLogEvent event) {
    try {
        process(event);
    } catch (Exception e) {
        log.warn("bad message, skipping: {}", e.getMessage()); // <- the loss is right here
    }
}

This looks defensive. It is the opposite. With the default acknowledgement mode, the container commits the offset once the listener method returns normally, and a caught-and-swallowed exception returns normally. So the moment the warning is logged and the code falls out of the catch, the container marks that message as done and advances the offset. The record is gone. Not dead-lettered, not retried, not visible anywhere except a log line that nobody reads until the customer asks where their data went.

The seductive part is that it never looks like an outage. The consumer keeps running, lag stays low, every dashboard is green. It has converted a loud failure into a silent one, which is strictly worse, because the system is now lying about its own completeness. The only honest options when a message cannot be processed are to keep retrying it or to set it aside somewhere durable. Swallowing it is neither.


Manual ack: the consumer chooses the commit point

log0's normalization consumer takes the commit point away from the container and owns it explicitly. The container factory is configured for manual acknowledgement:

java
factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL);

That one line changes the contract. The offset no longer commits when the method returns; it commits only when the code calls ack.acknowledge(). Now the question "is this message done?" has exactly one answer, and the code gets to decide it. The rule log0 picks is: a message is done when it is durable somewhere. There are two ways for that to be true, and the handler has a branch for each.

Three process tracks for the normalization consumer under AckMode.MANUAL. Track one, a clean log: normalize, publish to normalized-logs, save to ClickHouse, then ACK; durable because it is processed. Track two, a poison message that is quarantined: it throws, the catch wraps it in a DlqEvent, publishes it to raw-logs-dlq, then ACK; durable because it is in the dead-letter topic. Track three, the anti-pattern log0 refuses, drawn dashed in red: it throws, the catch logs and swallows it, the method returns normally so the container auto-commits, and the message is LOST, committed but never stored. The rule across all tracks: the offset commits only after the event is durable, and ack.acknowledge() is the last statement on every real pathThree process tracks for the normalization consumer under AckMode.MANUAL. Track one, a clean log: normalize, publish to normalized-logs, save to ClickHouse, then ACK; durable because it is processed. Track two, a poison message that is quarantined: it throws, the catch wraps it in a DlqEvent, publishes it to raw-logs-dlq, then ACK; durable because it is in the dead-letter topic. Track three, the anti-pattern log0 refuses, drawn dashed in red: it throws, the catch logs and swallows it, the method returns normally so the container auto-commits, and the message is LOST, committed but never stored. The rule across all tracks: the offset commits only after the event is durable, and ack.acknowledge() is the last statement on every real path

The real handler is small enough to read whole, and the shape of it is the whole point:

java
@KafkaListener(topics = "raw-logs", containerFactory = "kafkaListenerContainerFactory")
public void consume(RawLogEvent event, Acknowledgment ack) {
    try {
        NormalizedLogEvent normalized = normalizer.normalize(event);
        producer.publish(normalized);          // -> normalized-logs
        logEventRepository.save(normalized);    // -> ClickHouse
        ack.acknowledge();                      // durable: processed
    } catch (Exception e) {
        log.error("Error processing raw log message: {}", e.getMessage(), e);

        DlqEvent dlqEvent = DlqEvent.builder()
                .originalEvent(event)
                .errorMessage(e.getMessage())
                .failedAt("normalization-service")
                .failedAtTs(Instant.now())
                .build();

        dlqProducer.publish(event.getEventId(), dlqEvent); // -> raw-logs-dlq
        ack.acknowledge();                                 // durable: quarantined
    }
}

Notice where ack.acknowledge() sits: it is the last statement in both branches, and it is never reached until the message has gone somewhere it will survive a restart. On the happy path that is after the event is published downstream and written to ClickHouse. On the failure path it is after the original event has been wrapped and sent to the dead-letter topic. The offset advances in both cases, which is what stops a single poisoned message from stalling the partition behind it, but it advances only once the message is accounted for. A bad record is not skipped; it is moved.

The contrast with the bare-catch version is the entire lesson. Both versions catch the exception. Both keep the consumer alive. The difference is that one of them does something durable with the message before letting the offset move, and the other lets the offset move over a message it threw away.


What the dead-letter topic carries

A dead-letter queue is only useful if the thing put in it is enough to understand and recover the failure later. A raw payload with no context is barely better than a log line. So the DLQ record is an envelope, not only the original bytes:

java
DlqEvent.builder()
    .originalEvent(event)               // the full original payload, untouched
    .errorMessage(e.getMessage())       // why it failed
    .failedAt("normalization-service")  // which stage failed it
    .failedAtTs(Instant.now())          // when
    .build();

It carries the complete original event so nothing about the input is lost, plus the three facts needed when coming back to it: what went wrong, where, and when. The producer sends it to raw-logs-dlq keyed by the original event id:

java
kafkaTemplate.send(KafkaTopics.RAW_LOGS_DLQ, key, event); // key = original eventId

Keying by the original event id keeps the failure correlatable back to its source record and gives a future replay tool a stable handle. The intent is that raw-logs-dlq is a quarantine and a receipt, not a trash can: the message is out of the hot path so it cannot block anything, but it is fully preserved so it can be inspected, counted, and one day reprocessed. The honest status of "one day" is in the last section, because the topic is built and written correctly, but nothing drains it yet.


Proving it: inject poison and watch the split

A design like this is only worth the words if it is tested under actual failures, so I built a way to fail messages on purpose. The consumer has an environment-gated fault marker that is completely inert when unset:

java
@Value("${fault.inject-marker:}")
private String faultInjectMarker;
// ... inside consume(), before normalize():
if (!faultInjectMarker.isBlank() && event.getMessage() != null
        && event.getMessage().contains(faultInjectMarker)) {
    throw new IllegalStateException("fault-injection: poisoned message routed to DLQ");
}

When fault.inject-marker is empty, which is its default, this code does nothing and the consumer behaves exactly as in production. When it is set to a token, any log whose message contains that token throws on the way in, which exercises the real DLQ routing path without my having to take a dependency down or craft genuinely corrupt bytes. It is a deliberately small, off-by-default hook so the failure path is testable instead of theoretical.

Then I drove 600 logs through it: 300 carrying the poison marker, 300 clean. The result is the split this whole post is about.

A fault-injection run through the pipeline. 600 logs are sent, 300 poison plus 300 clean. The ingestion-gateway returns 600 times 202 with 0 errors and forwards all 600 to normalization. The 300 clean logs continue to incident creation, producing 1 incident at the clustering threshold of 10. The 300 poison logs branch off normalization down to raw-logs-dlq, the dead-letter topic, which receives all 300. The summary line: 0 messages dropped, 0 non-202 responses, and the consumer never blocks on a bad payloadA fault-injection run through the pipeline. 600 logs are sent, 300 poison plus 300 clean. The ingestion-gateway returns 600 times 202 with 0 errors and forwards all 600 to normalization. The 300 clean logs continue to incident creation, producing 1 incident at the clustering threshold of 10. The 300 poison logs branch off normalization down to raw-logs-dlq, the dead-letter topic, which receives all 300. The summary line: 0 messages dropped, 0 non-202 responses, and the consumer never blocks on a bad payload

Read the numbers off it. All 600 got a 202 at the gateway, because accept-fast does not know or care whether a payload will later parse. All 600 reached normalization. The 300 poison messages routed cleanly into raw-logs-dlq, all of them, none dropped. The 300 clean messages flowed straight through and still produced their incident, which matters more than it looks: it proves the poison sitting on the same topic and the same consumer group did not stall the clean records behind it. Because the consumer always acks, a bad message is set aside in milliseconds and the next message is processed immediately. The poison neither vanished nor jammed the line. Zero dropped, zero non-202, zero partition stalls. That is what "at-least-once" looks like when the promise is kept.


What is not done

The mechanism is sound and the gaps around it are real, so here is the honest list rather than a victory lap.

  • The DLQ send is fire-and-forget, and the ack does not wait for it. Look again: dlqProducer.publish(...) hands the record to the broker and returns, then ack.acknowledge() runs immediately, without confirming the DLQ write was durable. If the broker were failing at that exact moment, the original offset could be committed while the dead-letter write never lands, which is the one path back to silent loss. The correct version awaits the DLQ send's acknowledgement before acking the source. This is the most important fix on the list.
  • Nothing drains raw-logs-dlq yet. The envelope is designed for replay and is keyed for it, but there is no consumer that inspects, retries, or reprocesses dead-lettered messages today. They accumulate. A real system needs a DLQ reader, a triage view, and a guarded replay path back onto raw-logs.
  • There is no alert on DLQ depth. A message landing in the dead-letter topic is invisible to an operator right now. Quarantine without notification is only a quieter place to lose track of things. The thing to alert on is raw-logs-dlq receiving anything at all, because in steady state it should be empty.
  • The catch is too broad to tell transient from permanent. catch (Exception e) treats a genuinely unparseable payload and a momentary ClickHouse blip identically: both go straight to the DLQ. A transient downstream failure should be retried with backoff, not quarantined, otherwise a brief outage dead-letters a flood of perfectly good messages. Classifying the failure before routing it is unbuilt.
  • Acking last buys durability at the cost of duplicates. If the process crashes between save() and ack.acknowledge(), the offset never commits, the message is redelivered on restart, and it is normalized and written again. That is at-least-once by design: log0 chooses possible duplicates over possible loss. Deduplication downstream is real for incidents via the fingerprint, but the ClickHouse log_events write is not idempotent, so a redelivery can double a row.
  • All numbers are single-node (Docker Desktop, 512 MB per service, one Redpanda node, one ClickHouse node, k6). The 600/300/300 split is a correctness result, not a throughput one, and the routing behavior is a property of the code rather than the hardware.

The reason to lead with the fire-and-forget ack gap is that it is the kind of bug this whole post is about: the code looks like it cannot lose data, and there is still one window where it can. "Commit the offset last" is the right rule. "Commit the offset only after the durable write is confirmed" is the rule stated completely, and the gap between those two sentences is exactly where data goes missing.


Next: post 10, the incident lifecycle as an explicit state machine.. A poisoned message gets quarantined and a clean one becomes a log event, but a log event is not yet an incident. The next post follows what happens after detection: how a fingerprint becomes an open incident, how it moves through acknowledged and resolved without skipping steps, and why the transitions live in code as a state machine instead of a nullable status column that any update can scribble on.


Try log0

log0 is the platform this series is built on, an open, multi-tenant incident pipeline you can run yourself or use hosted.

Written by Ashmit JaiSarita Gupta. Find me on LinkedIn, GitHub, and X, and read the rest of the series on Hashnode.

Ashmit JaiSarita Gupta

Full-stack Software Engineer and the builder of log0. I write about backend systems, distributed systems, and the physics-flavored corners of engineering.

← Back to all posts

Turn log chaos into incident clarity

Get started
log0© 2026 log0, Inc.