Kafka with Java: Producers, Consumers, Delivery Semantics & Idempotency

Two weeks after the checkout API from Building a Complete Spring Boot API went live, we decoupled it. Instead of the checkout service calling the warehouse and notification services directly, it published an OrderPlaced event to a Kafka topic called order-events, and each downstream service consumed it at its own pace. Clean architecture. Fewer 2 a.m. pages. Then, on a busy Friday, a warehouse consumer pod got OOM-killed in the middle of a batch. Kubernetes restarted it. The consumer re-read the batch it had been processing — and printed every pick ticket twice. Two pallets went to one customer, the carrier billed us twice for the freight, and support spent the weekend apologizing.

Nothing was misconfigured. Kafka had done exactly what it promises: at-least-once delivery. The bug was ours — the consumer had no defense against seeing the same event twice. This post is that defense, end to end: how the producer guarantees no event is lost or duplicated on the way in, how the consumer owns its offsets, what the delivery-semantics table really means, and why "exactly-once" has a scope attached to it that most tutorials skip.

Where Kafka fits in the checkout story

After an order is placed, an OrderPlaced event flows through the order-events topic. The checkout service only produces; the warehouse, notifications, and loyalty services only consume. None of them know about each other. That decoupling is the whole point — but it moves the reliability contract from "did the HTTP call succeed?" to "what does this consumer do when it sees the same event twice?" Here's the shape of it:

Producer OrderEvent Producer idempotent, acks=all Topic: order-events durable log, 3 partitions P0 ord-1001, ord-1004, … ordered within partition P1 ord-1002, ord-1005, … P2 ord-1003, ord-1006, … Consumer group checkout-group one partition → one consumer C1 ← P0 C2 ← P1 C3 ← P2 each commits its own partition's offsets rebalance reassigns partitions on failure Key = orderId: one order's events always land on the same partition, in order. A 4th consumer would idle — partitions are the parallelism ceiling.

Three facts fall out of this picture that everything else builds on. A topic is a durable, append-only log, split into partitions — each partition is strictly ordered, but there is no ordering across partitions. A consumer group splits the partitions among its consumers, so the group scales by adding partitions, not consumers. And a record key decides the partition: key every event by the order id and one order's events always land on the same partition, in order.

The producer: acks, retries, and the idempotent producer

A producer's job is simple to state and easy to get wrong: every event reaches the topic, exactly once, even when the network misbehaves. Two settings control the tradeoff between speed and certainty. First, acks — how many brokers must confirm the write before the producer moves on:

acksMeaningFailure mode
0Fire and forget — no confirmation at allAny broker hiccup silently eats the event
1The partition leader confirmsLeader crashes before replicating → the event is gone, but you were told it succeeded
allEvery in-sync replica confirmsSurvives a leader crash; costs one replication round-trip of latency

Second, retries. Networks fail, leaders move, and a send that fails today may succeed on retry. But a naive retry has a trap: the producer sends a batch, the broker stores it, the acknowledgment is lost — so the producer retries and the broker stores the batch a second time. Retries without idempotence turn one network blip into duplicate events. That is exactly what the idempotent producer fixes: the broker tracks each producer by (producer-id, epoch, sequence-number) and discards any retried send it has already stored. A retry can no longer create a duplicate inside Kafka.

Here is the event and the producer. The config choices are the point of the listing:

package com.javamakeuse.orders;

/**
 * The event our checkout service publishes after an order is placed.
 * It flows through the "order-events" topic; warehouse, notifications
 * and loyalty services each consume it independently.
 */
public record OrderPlacedEvent(String orderId, String customerEmail, long totalCents) {

    public String toJson() {
        return "{\"orderId\":\"" + orderId + "\",\"customerEmail\":\"" + customerEmail
            + "\",\"totalCents\":" + totalCents + "}";
    }

    public static OrderPlacedEvent fromJson(String json) {
        String id = json.split("\"orderId\":\"")[1].split("\"")[0];
        String email = json.split("\"customerEmail\":\"")[1].split("\"")[0];
        long cents = Long.parseLong(json.split("\"totalCents\":")[1].replaceAll("[^0-9]", ""));
        return new OrderPlacedEvent(id, email, cents);
    }
}
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.clients.producer.RecordMetadata;
import org.apache.kafka.common.serialization.StringSerializer;

/**
 * The three settings that matter most on a producer:
 *  - acks=all: the leader waits until every in-sync replica has the record.
 *    Slower than acks=1, but a leader crash can no longer eat your event.
 *  - enable.idempotence=true: the broker dedups retried sends by
 *    (producer id, epoch, sequence number), so a retried send is never
 *    stored twice. The client enforces acks=all, retries>0 and
 *    max.in.flight.requests.per.connection<=5 with it (ConfigException otherwise).
 *  - linger.ms / batch.size: throughput knobs. Batching trades a few
 *    milliseconds of latency for far fewer broker round-trips.
 */
 /**
 * Publishes OrderPlaced events with the idempotent producer.
 *
 * The three settings that matter most on a producer:
 *  - acks=all: the leader waits until every in-sync replica has the record.
 *    Slower than acks=1, but a leader crash can no longer eat your event.
 *  - enable.idempotence=true: the broker dedups retried sends by
 *    (producer id, epoch, sequence number), so a retried send is never
 *    stored twice. The client enforces acks=all, retries>0 and
 *    max.in.flight.requests.per.connection<=5 with it (ConfigException otherwise).
 *  - linger.ms / batch.size: throughput knobs. Batching trades a few
 *    milliseconds of latency for far fewer broker round-trips.
 */
public class OrderEventProducer {

    private final Producer<String, String> producer;
    private final String topic;

    public OrderEventProducer(String bootstrapServers, String topic) {
        this.topic = topic;
        Properties props = new Properties();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());

        props.put(ProducerConfig.ACKS_CONFIG, "all");
        props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true");
        props.put(ProducerConfig.RETRIES_CONFIG, Integer.toString(Integer.MAX_VALUE));

        props.put(ProducerConfig.LINGER_MS_CONFIG, "20");
        props.put(ProducerConfig.BATCH_SIZE_CONFIG, Integer.toString(32 * 1024));

        this.producer = new KafkaProducer<>(props);
    }

    /** Package-visible for tests: inject a MockProducer here. */
    OrderEventProducer(Producer<String, String> producer, String topic) {
        this.producer = producer;
        this.topic = topic;
    }

    public void publish(OrderPlacedEvent event) {
        // Key = order id: every event for one order lands on the same
        // partition, so a consumer always sees them in order.
        ProducerRecord<String, String> record =
            new ProducerRecord<>(topic, event.orderId(), event.toJson());
        producer.send(record, (RecordMetadata meta, Exception ex) -> {
            if (ex != null) {
                // Production code: alert, then route to a retry topic or DLQ.
                // Swallowing this silently is how events vanish without a trace.
                throw new IllegalStateException("Failed to publish " + event.orderId(), ex);
            }
        });
    }

    public void close() {
        producer.close();
    }
}

Two details in that listing are worth pausing on, because I verified them against the actual client (version 3.9.0) rather than trusting memory. With enable.idempotence=true, the client fills in acks=all (reported as -1), retries=2147483647 (Integer.MAX_VALUE), and max.in.flight.requests.per.connection=5 — and it rejects anything weaker with a ConfigException: acks=1 fails with "Must set acks to all in order to use the idempotent producer. Otherwise we cannot guarantee idempotence.", as do retries=0 and max.in.flight.requests.per.connection=10. Idempotence is not a suggestion the broker politely considers; the client refuses to run without the settings that make it true. (The linger.ms=20 / batch.size=32KB pair just batches small sends for a few milliseconds — pure throughput tuning, no correctness impact.)

Decision rule: default every producer to acks=all with the idempotent producer on, and only relax it when you can name the latency budget that forces you to. The replication round-trip is a bounded, measurable cost; duplicate or lost events are an unbounded debugging cost.

The consumer: the poll loop, offsets, and who owns the commit

A consumer does not get events pushed to it. It polls — asks the broker "what's new since the offset I last committed?" — processes the batch, then tells the broker how far it got. That "how far I got" is the offset commit, and who moves it, and when, decides your delivery semantics on the read side:

Commit strategyWhen the offset movesFailure mode
Auto-commit (the default, every 5000 ms)A background thread commits whatever was last polledCrash after polling but before processing → those records are never re-read: lost events
commitSync() after the batchAfter your code finished processingCrash mid-batch → the batch is replayed: duplicates, but nothing lost
commitAsync()Fire-and-forget, no error visibilityA failed commit is silent; commit ordering across calls isn't guaranteed

The auto-commit default (enable.auto.commit=true, interval 5000 ms — both verified against the client) is the trap: it commits offsets for records you have polled, not records you have processed. A crash in that gap loses events permanently, with no error and no trace. Any consumer with side effects — printing pick tickets, charging cards, crediting loyalty points — must own its commits.

Here is the warehouse consumer, written the safe way: auto-commit off, process the batch, then commitSync():

package com.javamakeuse.orders;

import java.time.Duration;
import java.util.Collection;
import java.util.List;
import java.util.Properties;

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRebalanceListener;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.serialization.StringDeserializer;

/**
 * At-least-once consumer for the "order-events" topic.
 *
 * The two load-bearing decisions:
 *  1. enable.auto.commit=false — WE own the commit. Auto-commit would
 *     advance offsets for records we have not finished processing.
 *  2. process-then-commit — a crash replays at most the uncommitted tail,
 *     and the idempotent OrderStore makes that replay a no-op.
 */
public class OrderEventConsumer {

    /** The downstream write, made idempotent by keying on the order id. */
    public interface OrderStore {
        /** @return true if this order id is seen for the first time */
        boolean markProcessedIfNew(String orderId);
        void apply(OrderPlacedEvent event);
    }

    private final KafkaConsumer<String, String> consumer;
    private final OrderStore store;

    public OrderEventConsumer(String bootstrapServers, String groupId, OrderStore store) {
        this.store = store;
        Properties props = new Properties();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());

        props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
        // Small batches bound the replay after a crash.
        props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, "100");

        this.consumer = new KafkaConsumer<>(props);
    }

    public void run(String topic) {
        consumer.subscribe(List.of(topic), new ConsumerRebalanceListener() {
            @Override
            public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
                // About to lose these partitions: commit finished work first,
                // so the new owner does not replay what we already applied.
                consumer.commitSync();
            }

            @Override
            public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
                // Offsets resume from the last commit; nothing to do here.
            }
        });

        while (true) {
            ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
            for (ConsumerRecord<String, String> record : records) {
                OrderPlacedEvent event = OrderPlacedEvent.fromJson(record.value());
                if (store.markProcessedIfNew(event.orderId())) {
                    store.apply(event);   // downstream DB write, inventory, ...
                }                         // else: redelivery — already applied, skip.
            }
            if (!records.isEmpty()) {
                consumer.commitSync();
            }
        }
    }
}

A word on the rebalance listener, because the opening incident was a rebalance story too. When a consumer joins, leaves, or dies, the group rebalances: partitions are revoked from some consumers and assigned to others. If your consumer dies mid-batch without committing, the new owner of that partition starts from the last committed offset — replaying everything since. Committing in onPartitionsRevoked shrinks that replay window: before you hand a partition back, you checkpoint the work you finished. It doesn't eliminate replays (a hard crash can't run the listener), which is exactly why the consumer also needs to be idempotent — the next section.

Decision rule: disable auto-commit on every consumer that has side effects. Commit after processing, never before — and treat "commit before processing" as a data-loss bug no matter how convenient it looks.

Delivery semantics: pick your guarantee, know its price

Everything so far fits in one table. Read it as a menu where each row costs something:

GuaranteeProducer sideConsumer sideWhat can still go wrong
At-most-onceacks=0, no retriesCommit before (or regardless of) processing — e.g. auto-commitEvents are silently lost on any failure
At-least-onceacks=all + retriesProcess, then commit (manual)Crashes and rebalances cause duplicates
Effectively-onceIdempotent producerAt-least-once + idempotent consumer (dedup by key)Duplicates are delivered but absorbed; the dedup store must be durable
Exactly-once in the brokerTransactional producerisolation.level=read_committedSee the warning below — the broker's transaction does not cover your database

Notice the shape of the answer: Kafka gives you at-least-once delivery and idempotent producing for free, and the consumer is where "once" actually happens — by making the downstream write safe to repeat. That is the effectively-once row, and it is the row real systems live on.

The exactly-once warning: the broker's transaction is not your transaction

Kafka does offer exactly-once inside the broker, via transactions. A transactional producer can atomically commit two things together: the consumer offsets it read up to, and the records its processing produced. Either both become visible or neither does — a downstream consumer reading with isolation.level=read_committed never sees half a processing step. The transactional.id also lets the broker fence zombie producers: after a restart, the old instance's in-flight transaction is aborted, so two instances can never both write.

package com.javamakeuse.orders;

import java.util.Map;
import java.util.Properties;

import org.apache.kafka.clients.consumer.ConsumerGroupMetadata;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.Producer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringSerializer;

/**
 * Broker-side exactly-once: the transaction atomically commits
 * (a) the consumer offsets you consumed up to, and
 * (b) the records your processing produced.
 *
 * If the transaction aborts, neither the offsets nor the records become
 * visible — so a downstream consumer of this topic never sees half of
 * a processing step.
 *
 * WARNING (read the post): this is exactly-once *inside Kafka*.
 * A JDBC write to your own database is not part of this transaction —
 * making the whole pipeline exactly-once still needs an idempotent
 * consumer on the receiving side.
 */
public class TransactionalOrderProducer {

    private final Producer<String, String> producer;

    public TransactionalOrderProducer(String bootstrapServers, String transactionalId) {
        Properties props = new Properties();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers);
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true");
        // One stable id per producer instance: the broker uses it to fence
        // zombie producers after a restart (the old instance's writes abort).
        props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, transactionalId);
        this.producer = new KafkaProducer<>(props);
        this.producer.initTransactions();
    }

    /**
     * Atomically: mark the consumed offsets AND emit the result event.
     * Consumers reading with isolation.level=read_committed only ever see
     * committed transactions — never the half-written middle.
     */
    public void processAndEmit(ConsumerGroupMetadata group,
                               Map<TopicPartition, OffsetAndMetadata> offsets,
                               OrderPlacedEvent result) {
        producer.beginTransaction();
        try {
            producer.send(new ProducerRecord<>("order-events", result.orderId(), result.toJson()));
            producer.sendOffsetsToTransaction(offsets, group);
            producer.commitTransaction();
        } catch (Exception e) {
            producer.abortTransaction();
            throw e;
        }
    }
}

Now the warning, stated plainly: broker exactly-once is not end-to-end exactly-once. The transaction above covers offsets and produced records — both inside Kafka. It does not cover the JDBC write your consumer made to its own database a moment earlier. If the process crashes between "order row committed to Postgres" and "Kafka transaction committed", the consumer restarts, replays the record, and writes the order row a second time. Kafka kept its promise; your database still got a duplicate. There is no distributed transaction spanning Kafka and your database here — and that is a feature, not a gap: two-phase commit across systems is how you get outages that take down both sides at once.

The honest vocabulary is: "exactly-once inside Kafka via transactions" or "effectively-once end-to-end via an idempotent consumer" — never unscoped "exactly-once". If the database-transaction half of this story is what you want next — what @Transactional actually guarantees, the proxy behind it, and isolation levels — that is Transactions: @Transactional, Isolation & Propagation. Kafka transactions and database transactions are two different machines that happen to share a name; don't let the shared word fool you into thinking one covers the other.

Decision rule: say "exactly-once" only with a scope attached. Unscoped "exactly-once" is how the duplicate-shipment bug gets reintroduced by the next well-meaning refactor.

The idempotent consumer: making redelivery harmless

If delivery is at-least-once, duplicates are not an edge case — they are the contract. The fix is not to prevent redelivery (you can't, not across crashes and rebalances); it is to make processing safe to repeat. An idempotent consumer keeps a record of what it already applied — a processed_order_ids table, a unique constraint on the order id, a Redis set — and skips anything it has seen. The downstream write becomes: apply if new, ignore if seen.

The demo below proves the whole chain with no broker, using the MockProducer and MockConsumer that ship inside the kafka-clients jar. It (a) produces three OrderPlaced events with the idempotent-producer config from the previous section, (b) consumes them with at-least-once semantics — process, then commit — plus a dedup store keyed by order id, and (c) simulates a redelivery (the batch replays, as after a rebalance) to show the dedup store absorbing it:

package com.javamakeuse.orders;

import java.time.Duration;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Properties;
import java.util.Set;

import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerGroupMetadata;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.clients.consumer.OffsetResetStrategy;
import org.apache.kafka.clients.producer.MockProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.serialization.StringSerializer;

/**
 * End-to-end demo with NO broker, using the mock clients that ship in the
 * kafka-clients jar:
 *  (a) produce OrderPlaced events with the idempotent-producer config,
 *  (b) consume them with at-least-once semantics + an idempotent consumer
 *      that dedups by order id,
 *  (c) simulate a redelivery and prove the dedup store absorbs it.
 */
public class OrderEventDemo {

    static final String TOPIC = "order-events";

    public static void main(String[] args) {
        // ---- 1. The idempotent producer's config, verified against the real client ----
        Properties props = new Properties();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true");
        ProducerConfig cfg = new ProducerConfig(props);
        System.out.println("[config] enable.idempotence=true  =>  acks="
            + cfg.getString(ProducerConfig.ACKS_CONFIG) + " (all), retries="
            + cfg.getInt(ProducerConfig.RETRIES_CONFIG) + ", max.in.flight="
            + cfg.getInt(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION));

        // Consumer defaults, straight from the client:
        ConsumerConfig ccfg = new ConsumerConfig(Map.of(
            ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092",
            ConsumerConfig.GROUP_ID_CONFIG, "checkout-group",
            ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,
                "org.apache.kafka.common.serialization.StringDeserializer",
            ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,
                "org.apache.kafka.common.serialization.StringDeserializer"));
        System.out.println("[config] consumer defaults: enable.auto.commit="
            + ccfg.getBoolean(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG)
            + ", auto.commit.interval.ms=" + ccfg.getInt(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG)
            + ", max.poll.records=" + ccfg.getInt(ConsumerConfig.MAX_POLL_RECORDS_CONFIG)
            + ", isolation.level=" + ccfg.getString(ConsumerConfig.ISOLATION_LEVEL_CONFIG));

        // ---- 2. Produce three OrderPlaced events ----
        MockProducer<String, String> producer =
            new MockProducer<>(true, new StringSerializer(), new StringSerializer());
        OrderPlacedEvent[] orders = {
            new OrderPlacedEvent("ord-1001", "mia@example.com", 14999),
            new OrderPlacedEvent("ord-1002", "raj@example.com", 7999),
            new OrderPlacedEvent("ord-1003", "ana@example.com", 21999),
        };
        for (OrderPlacedEvent o : orders) {
            producer.send(new ProducerRecord<>(TOPIC, o.orderId(), o.toJson()));
        }
        System.out.println("[produce] " + producer.history().size()
            + " OrderPlaced events sent to '" + TOPIC + "'");

        // ---- 3. Consume with at-least-once: process, then commit ----
        MockConsumer<String, String> consumer = new MockConsumer<>(OffsetResetStrategy.EARLIEST);
        TopicPartition tp = new TopicPartition(TOPIC, 0);
        consumer.assign(List.of(tp));
        consumer.updateBeginningOffsets(Map.of(tp, 0L));
        enqueue(consumer, producer.history());

        Set<String> appliedOrderIds = new HashSet<>();   // the downstream "orders" table
        int[] first = consumeOnce(consumer, appliedOrderIds);
        consumer.commitSync();
        System.out.println("[consume] first pass: processed=" + first[0]
            + ", duplicates-skipped=" + first[1]
            + ", committed-offset=" + consumer.committed(tp).offset()
            + ", orders-in-store=" + appliedOrderIds.size());

        // ---- 4. Simulate a redelivery: the broker replays records we already applied ----
        consumer.seek(tp, 0L);                            // e.g. after a rebalance or lost commit
        enqueue(consumer, producer.history());            // broker re-delivers the batch
        int[] second = consumeOnce(consumer, appliedOrderIds);
        consumer.commitSync();
        System.out.println("[consume] after redelivery: processed=" + second[0]
            + ", duplicates-skipped=" + second[1]
            + ", committed-offset=" + consumer.committed(tp).offset()
            + ", orders-in-store=" + appliedOrderIds.size() + " (dedup held)");

        // ---- 5. Broker-side exactly-once: consume-position + produce, one transaction ----
        MockProducer<String, String> txProducer =
            new MockProducer<>(true, new StringSerializer(), new StringSerializer());
        txProducer.initTransactions();
        txProducer.beginTransaction();
        txProducer.send(new ProducerRecord<>(TOPIC, "ord-1004", "{\"orderId\":\"ord-1004\"}"));
        txProducer.sendOffsetsToTransaction(
            Map.of(tp, new OffsetAndMetadata(3L)), new ConsumerGroupMetadata("checkout-group"));
        txProducer.commitTransaction();
        System.out.println("[txn] consume-offsets + produced record committed atomically");
    }

    private static void enqueue(MockConsumer<String, String> consumer,
                                List<ProducerRecord<String, String>> history) {
        long offset = 0;
        for (ProducerRecord<String, String> r : history) {
            consumer.addRecord(new ConsumerRecord<>(TOPIC, 0, offset++, r.key(), r.value()));
        }
    }

    /** Idempotent consumer: dedups by order id. Returns {processed, duplicatesSkipped}. */
    private static int[] consumeOnce(MockConsumer<String, String> consumer,
                                     Set<String> appliedOrderIds) {
        int processed = 0, skipped = 0;
        ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
        for (ConsumerRecord<String, String> record : records) {
            OrderPlacedEvent event = OrderPlacedEvent.fromJson(record.value());
            if (appliedOrderIds.add(event.orderId())) {
                processed++;        // first sighting: apply the downstream write
            } else {
                skipped++;          // redelivery: already applied, skip it
            }
        }
        return new int[] { processed, skipped };
    }
}

The genuine console output, run against kafka-clients 3.9.0 with no broker anywhere:

[config] enable.idempotence=true  =>  acks=-1 (all), retries=2147483647, max.in.flight=5
[config] consumer defaults: enable.auto.commit=true, auto.commit.interval.ms=5000, max.poll.records=500, isolation.level=read_uncommitted
[produce] 3 OrderPlaced events sent to 'order-events'
[consume] first pass: processed=3, duplicates-skipped=0, committed-offset=3, orders-in-store=3
[consume] after redelivery: processed=0, duplicates-skipped=3, committed-offset=3, orders-in-store=3 (dedup held)
[txn] consume-offsets + produced record committed atomically

Read the two [consume] lines as the whole post in miniature. First pass: 3 processed, offset committed at 3. After the redelivery: 0 processed, 3 skipped, the store still holds exactly 3 orders. The delivery was at-least-once; the outcome was effectively-once. Note also the [config] lines: the client's real defaults — auto-commit on with a 5-second interval — are exactly why the manual-commit consumer above sets it to false explicitly.

1. Normal batch poll → 3 records ord-1001 … ord-1003 apply ×3 orders-in-store: 3 commit offset 3 2. Redelivery broker replays same 3 records dedup store all order ids seen skip ×3 — nothing applied twice 3. Without the dedup store, step 2 re-applies all 3 orders — double shipment, double loyalty points, one angry customer. Delivery was at-least-once. The outcome was effectively-once — because the consumer was idempotent.

One production caveat the demo's HashSet hides: the dedup store must survive restarts. An in-memory set forgets everything when the pod dies — which is precisely when redelivery arrives. In production this is a database table with a unique constraint on the order id (the insert itself is the atomic "mark processed" — a duplicate insert fails and you skip), or an external store with a TTL matched to your retention. The principle from the demo still holds; only the storage changes.

Decision rule: key every event by the entity it describes, and make the consumer idempotent on that key in durable storage. Redelivery is a when, not an if — design the consumer for the replay on day one, not after the first incident.

When processing fails: retry topics and the dead-letter queue

Idempotency handles duplicates. It doesn't handle a record that can never be processed — the poison pill: malformed JSON, an order id that violates a constraint, a downstream service that rejects it every time. Retrying it in place blocks the whole partition behind it (Kafka only moves forward per partition), so the standard pattern moves the failure sideways instead:

  • Retry topics — order-events-retry-1, order-events-retry-2, …: a failed record is re-published here with a backoff, and a consumer retries it up to N times. Each attempt is a fresh poll, so the main partition keeps flowing.
  • Dead-letter queue (DLQ) — order-events-dlq: after the attempts are exhausted, the record lands here with its error attached. Nothing auto-processes the DLQ; a human (or a fixed consumer) investigates and replays it deliberately.

The router below is the whole pattern in miniature — the same code ran in the test at the end:

package com.javamakeuse.orders;

import java.time.Duration;
import java.util.List;
import java.util.Map;

import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.MockConsumer;
import org.apache.kafka.clients.consumer.OffsetResetStrategy;
import org.apache.kafka.clients.producer.MockProducer;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.TopicPartition;
import org.apache.kafka.common.serialization.StringSerializer;

/**
 * Retry-with-backoff + dead-letter routing.
 *
 * A record that always fails must never park the partition forever:
 * after MAX_ATTEMPTS it leaves the retry chain and lands in the DLQ,
 * where a human (or a fixed consumer) can deal with it.
 */
public class RetryRouter {

    static final int MAX_ATTEMPTS = 3;

    private final MockProducer<String, String> router =
        new MockProducer<>(true, new StringSerializer(), new StringSerializer());

    /** Route one record: process it, or forward it down the retry chain. */
    public void handle(ConsumerRecord<String, String> record, int attempt) {
        try {
            process(record);                       // your business logic
        } catch (TransientException e) {
            if (attempt < MAX_ATTEMPTS) {
                router.send(new ProducerRecord<>(
                    "order-events-retry-" + (attempt + 1), record.key(), record.value()));
            } else {
                router.send(new ProducerRecord<>(
                    "order-events-dlq", record.key(), record.value()));
            }
        }
    }

    private void process(ConsumerRecord<String, String> record) {
        if (record.value().contains("POISON")) {
            throw new TransientException("simulated poison record");
        }
    }

    static class TransientException extends RuntimeException {
        TransientException(String msg) { super(msg); }
    }

    public List<ProducerRecord<String, String>> routed() {
        return router.history();
    }

    public static void main(String[] args) {
        RetryRouter router = new RetryRouter();
        MockConsumer<String, String> consumer = new MockConsumer<>(OffsetResetStrategy.EARLIEST);
        TopicPartition tp = new TopicPartition("order-events", 0);
        consumer.assign(List.of(tp));
        consumer.updateBeginningOffsets(Map.of(tp, 0L));
        // A poison record arrives on its 4th attempt: retries are exhausted.
        consumer.addRecord(new ConsumerRecord<>("order-events", 0, 0L, "ord-9001", "POISON"));
        ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
        for (ConsumerRecord<String, String> r : records) {
            router.handle(r, 3);
        }
        for (ProducerRecord<String, String> r : router.routed()) {
            System.out.println("poison record routed to: " + r.topic() + " (key=" + r.key() + ")");
        }
    }
}
poison record routed to: order-events-dlq (key=ord-9001)

The poison record left the main topic instead of wedging it. Two things to carry into production: the retry topics need their own consumers with real backoff (a fixed delay or exponential sleep before re-publishing — hammering a failing downstream in a tight loop is a self-inflicted DDoS), and the DLQ needs an owner and an alert, or it becomes a write-only graveyard nobody reads.

Decision rule: a record that fails forever must leave the main partition. Bounded retries, then the DLQ — never an infinite in-place retry loop, and never a silent catch-and-drop.

Testing this without a broker

Everything above was verified without a running Kafka, and your tests can work the same way: MockProducer captures everything send()ed in history() so you can assert on the records your code produced, and MockConsumer lets you addRecord() canned events, drive the poll loop, and assert on committed offsets. That covers logic — the dedup check, the routing, the commit strategy. For the parts mocks can't cover (serialization against a real broker, partition assignment, rebalance behavior), run an integration test against a real broker with Testcontainers' Kafka support, as covered in Testcontainers: Real Databases in Tests — the same container pattern, with a Kafka broker instead of a database.

What's next

Kafka is reliable now — but there is one gap left in the checkout story, on the producing side. The checkout service writes the order row to its database and publishes the OrderPlaced event as two separate steps: if the process dies between them, the database and the topic disagree forever. The next post in this track, " Transactional Outbox: DB + Kafka Without Losing Messages", closes that gap: how to publish the event as part of the database transaction itself, so the order row and the event can never disagree.

Field check before you move on: take the demo above and break it on purpose — comment out the dedup check in consumeOnce, rerun the redelivery, and watch orders-in-store jump to 6. Then restore the check and add a fourth record whose orderId collides with an existing one but carries a different totalCents: decide what your real store should do with it — first-write-wins? last-write-wins? reject and alert? — and encode that decision in a code comment. That comment is the idempotency policy your team will argue about in the next incident review; write it now, while it's cheap.

Continue: Java Learning Roadmap 2026

Comments

Popular posts from this blog

JSP Servlet Interview Questions For Freshers Series 1

Java Banking Finance Services and Insurance (BFSI) domain interview questions

Java program to check even or odd number