skip to content

As a principal engineer, how do you configure KafkaTemplate/ProducerFactory for durable, exactly-once-style delivery, and what are the trade-offs?

level: principalimportance: should knowfreq 40%

answer

  1. acks=all + min.insync.replicas>=2 = durable
  2. enable.idempotence -> per-partition effectively-once, dedupes retries
  3. transactionIdPrefix -> executeInTransaction / KafkaTransactionManager
  4. consumers read_committed; transactional.id fencing
  5. exactly-once is IN Kafka only -> outbox for DB

basics

~20 s

Set acks=all and enable.idempotence=true on the ProducerFactory for durable, non-duplicating writes. For atomic multi-message sends, add a transactionIdPrefix so the template runs Kafka transactions. The trade-off is lower throughput and higher latency for stronger guarantees.

solid answer

~40 s

Durability and de-duplication are set on the ProducerFactory's config: acks=all makes the leader wait for all in-sync replicas, and enable.idempotence=true (which also forces acks=all, capped in-flight, and unlimited-ish retries) prevents duplicate records from producer retries, giving effectively-once *per partition*. For true atomicity across multiple records or a consume-transform-produce loop, set a transactionIdPrefix on the DefaultKafkaProducerFactory; KafkaTemplate then supports executeInTransaction and integrates with KafkaTransactionManager, so a batch of sends commits or aborts together and consumers with read_committed skip aborted records. Trade-offs: transactions and acks=all add latency and reduce throughput, transactional.id fencing complicates scaling, and 'exactly-once' only holds within Kafka — external side effects (DB writes) need the transactional outbox or idempotent consumers. Also plan delivery.timeout.ms and error handling for the futures.

code

java · 32 lines
java
@Bean
public DefaultKafkaProducerFactory<String, OrderEvent> producerFactory() {
    Map<String, Object> props = new HashMap<>();
    props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "broker:9092");
    props.put(ProducerConfig.ACKS_CONFIG, "all");
    props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
    DefaultKafkaProducerFactory<String, OrderEvent> pf =
            new DefaultKafkaProducerFactory<>(props);
    pf.setTransactionIdPrefix("orders-tx-"); // enables transactions
    return pf;
}

@Bean
public KafkaTransactionManager<String, OrderEvent> ktm(
        ProducerFactory<String, OrderEvent> pf) {
    return new KafkaTransactionManager<>(pf);
}

@Service
class OrderPublisher {
    private final KafkaTemplate<String, OrderEvent> template;
    OrderPublisher(KafkaTemplate<String, OrderEvent> template) { this.template = template; }

    void publishAtomically(OrderEvent a, OrderEvent b) {
        // both sends commit together or both abort
        template.executeInTransaction(t -> {
            t.send("orders", a.id(), a);
            t.send("audit", a.id(), b);
            return null;
        });
    }
}

go deeper

for a junior

Know acks=all means wait for replicas and is more durable than acks=1/0.

for a middle

Explain acks levels and that enable.idempotence prevents duplicate records from retries.

for a senior

Configure idempotence + transactions, know read_committed and transactionIdPrefix, and the throughput trade-off.

for a principal

Reason about the exactly-once boundary, outbox for cross-system consistency, transactional.id fencing/scaling, and durability sizing (ISR/replication).

**The delivery-semantics spectrum.** Kafka producers can give *at-most-once*, *at-least-once*, or *effectively/exactly-once* depending on config. As a principal you choose deliberately. **1. Durability — `acks`.** Configured on the ProducerFactory (`ProducerConfig.ACKS_CONFIG`): - `acks=0` at-most-once, fire-and-forget, can lose data on broker failure. - `acks=1` leader ack only — a leader crash before replication loses the record. - `acks=all` (`-1`) leader waits for all *in-sync replicas* (ISRs). Combine with broker `min.insync.replicas>=2` and replication factor ≥3 so a single broker loss can't lose acknowledged data. This is the durability baseline. **2. No duplicates from retries — `enable.idempotence=true`.** Producer retries (needed for reliability) can otherwise write a record twice if an ack is lost. The **idempotent producer** tags records with a producer id + sequence number so the broker deduplicates retries *per partition*. Enabling it implicitly sets `acks=all`, bounds `max.in.flight.requests.per.connection<=5` while *preserving ordering*, and allows generous retries. Result: **effectively-once within a partition** with retries — the sensible default for most producers (and it's on by default in recent Kafka clients). **3. Atomic multi-record / read-process-write — transactions.** Idempotence dedupes single records but doesn't make *a group* of sends atomic, nor tie sends to consumer-offset commits. For that, set a **`transactionIdPrefix`** on `DefaultKafkaProducerFactory` (`spring.kafka.producer.transaction-id-prefix`). Then: - `KafkaTemplate` gains `executeInTransaction(t -> { t.send(...); t.send(...); })` — all sends commit or abort together. - With `@Transactional` + `KafkaTransactionManager`, a consume-transform-produce loop commits produced records **and** consumer offsets in one Kafka transaction — the basis of Kafka's exactly-once processing. - Consumers must set `isolation.level=read_committed` to skip records from aborted transactions. - Each transactional producer has a unique `transactional.id`; Kafka uses it to **fence** zombie/duplicate producers after failover. **Trade-offs to articulate.** - **Throughput/latency:** `acks=all` waits for replicas; transactions add begin/commit round-trips and hold partitions — measurable overhead. Tune `linger.ms`/`batch.size`/`compression.type` to claw back throughput. - **Scaling complexity:** `transactional.id` must be stable and unique per logical producer; naive horizontal scaling can cause fencing surprises. Spring manages a pool but you must reason about instance identity. - **Boundary of the guarantee:** exactly-once is *within Kafka*. A `send()` plus a database `INSERT` is **not** atomic across systems. To make the DB and Kafka consistent you use the **transactional outbox pattern** (write event to an outbox table in the DB transaction, relay to Kafka) or design **idempotent consumers** — don't claim end-to-end exactly-once from producer config alone. - **Error handling still required:** even with all this, observe the send futures, set `delivery.timeout.ms` sensibly, and have a dead-letter / retry strategy; durability config doesn't remove the need to handle a genuinely failed send. **Decision guide.** - Internal events, occasional dupes tolerable → `acks=all` + idempotence (default). Simple, fast enough. - Must not duplicate and must be atomic across several topics / with offset commits → transactions. - Must be consistent with a database → outbox pattern, not Kafka transactions spanning the DB. **Spring specifics.** All of this is producer-config on the ProducerFactory plus (for transactions) `transactionIdPrefix` and a `KafkaTransactionManager` bean. `KafkaTemplate.executeInTransaction` is the programmatic entry; `@Transactional` on a method whose transaction manager is the KafkaTransactionManager is the declarative one. Recovery/retry for listeners is a separate consumer concern (`DefaultErrorHandler`, dead-letter topics).

  • Does enabling Kafka transactions on the producer give you exactly-once across a Kafka send and a database write?
    No. Kafka transactions are exactly-once only within Kafka (atomic sends + offset commits). A DB write is a separate system and isn't part of the Kafka transaction. For consistency across both, use the transactional outbox pattern or make the consumer idempotent.
  • Why does enable.idempotence limit max.in.flight.requests.per.connection?
    The idempotent producer must preserve per-partition ordering while deduplicating retries using sequence numbers. To keep sequences correct under retries it caps in-flight requests (<=5), so a retried batch can't be reordered ahead of a later one.
  • What must consumers set to benefit from producer transactions?
    isolation.level=read_committed, so they skip records belonging to aborted transactions and only see committed data. With read_uncommitted (the default) they'd see aborted records too.

saying these in an interview costs you the question

  • Claiming Kafka transactions give exactly-once across Kafka and a database
  • Thinking idempotence alone makes multiple sends atomic
  • Ignoring min.insync.replicas / replication factor so acks=all still risks loss
  • Assuming read_committed is the consumer default (it is read_uncommitted)
  • Treating transactional.id as free to duplicate across scaled instances

context