Two producer questions come up again and again in Kafka interviews. The first is about correctness: “can a producer retry create a duplicate, and what does enable.idempotence do?” The second is about performance: “how would you increase producer throughput?”, which leads to batching, linger.ms and compression. They belong together because both live in the same machinery: the record accumulator that batches records per partition and the sender thread that ships those batches and retries them.
This guide explains how the producer turns send() calls into batches, how sequence numbers make retries safe, where idempotence stops protecting you, and how to tune for throughput. It targets Kafka 3.8/4.0 and the Java client. Java code and CLI commands are illustrative and were not run against a cluster; the deduplication model in the verification section is a Node.js script that was run.
Before you start
You should know that a topic has partitions with one leader broker each, and what acks=all means; Kafka replication, ISR and acks covers that. Basic Java, including callbacks and Future, is enough for the code.
The short answer
send() does not hit the network: it serializes the record, picks a partition and appends it to an in-memory batch for that partition. A background sender thread sends a batch when it reaches batch.size (16 KB) or has waited linger.ms (5 ms by default in Kafka 4.0, 0 before), compressing whole batches if compression.type is set. Because the network can lose an acknowledgement, a retry could write a batch twice; the idempotent producer, on by default since Kafka 3.0, prevents that by tagging batches with a producer ID and per-partition sequence numbers so the broker can drop duplicates and reject gaps. It only covers retries within one producer session, not an application that restarts and sends the same event again.
How it works
The path of a record through the Java producer:
- Serialize the key and value, then partition: explicit partition, hash of the key, or the sticky partitioner for keyless records.
- Accumulate: the record is appended to the open batch for its partition in the
RecordAccumulator, which is bounded bybuffer.memory(32 MB). If the buffer is full,send()blocks for up tomax.block.ms(60 seconds) and then throws. - Send: the sender thread collects ready batches, groups them by leader broker, and sends one produce request per broker, up to
max.request.size(1 MB). - Complete: when the response arrives, each record’s
Futureand callback complete with the partition and offset, or with an exception after retries are exhausted withindelivery.timeout.ms(2 minutes).
Retries create two risks. If the broker wrote a batch but the response was lost, resending writes it again: a duplicate. If up to 5 requests are in flight (max.in.flight.requests.per.connection) and an early one fails and is retried after a later one succeeded, records land out of order: a reorder.
Idempotence fixes both. At startup the producer gets a producer ID (PID) and epoch from the broker. Every batch carries the PID and the sequence number of its first record, counted per partition. The leader remembers the sequences of the last five batches for each PID and partition:
- A batch it has already written is acknowledged again without being appended.
- A batch whose sequence skips ahead is rejected with an out-of-order sequence error, so the producer resends the missing batch first.
That five-batch memory is why idempotence requires max.in.flight.requests.per.connection of 5 or less, along with acks=all and retries enabled. If you never set enable.idempotence and set a conflicting value such as acks=1, the client quietly disables idempotence; if you set it to true explicitly with a conflicting value, the producer refuses to start with a configuration error.
Step-by-step walkthrough
Step 1: Make idempotence explicit
Relying on the default is fragile, because a later edit like acks=1 silently removes it. Setting it explicitly turns that edit into a startup failure.
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 5);Step 2: Send asynchronously and handle the result in a callback
Calling .get() on every send() waits for a full round trip per record, so batches never fill and throughput collapses to a few thousand records per second. A callback keeps the pipeline full while still surfacing failures.
// Slow: one network round trip per record
producer.send(record).get();
// Fast: batching works, failures still handled
producer.send(record, (metadata, exception) -> {
if (exception != null) log.error("failed to send {}", record.key(), exception);
});Callbacks run on the sender thread, so keep them short: count, log or hand off, but do not block.
Step 3: Tune batching for throughput
Throughput comes from fewer, larger requests. Give batches room to grow and a little time to fill, then compress them.
props.put(ProducerConfig.BATCH_SIZE_CONFIG, 64 * 1024); // 64 KB per partition batch
props.put(ProducerConfig.LINGER_MS_CONFIG, 20); // wait up to 20 ms to fill
props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "zstd"); // or lz4 for lower CPUbatch.size is a per-partition upper bound, so memory use scales with the number of partitions being written. linger.ms adds latency only when traffic is light; under heavy load batches fill before the timer fires. Compression works on whole batches, so larger batches compress better and reduce network, disk and replication traffic. With the broker or topic compression.type left at its default of producer, brokers store the batches as the producer compressed them.
Step 4: Close the producer properly
Records waiting in the accumulator exist only in memory. The sender is a daemon thread, so if the JVM exits without flush() or close(), those records are lost silently.
try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
for (Event e : events) producer.send(new ProducerRecord<>("events", e.id(), e.json()), callback);
} // close() waits for in-flight and buffered recordsWorked scenario
A billing service exports invoice events. During a broker restart, the downstream ledger records some invoices twice. Investigation shows two separate causes.
First, the team had set acks=1 “for performance” a year earlier. Without an explicit enable.idempotence, that change disabled idempotence. When the leader moved, requests whose responses were lost were retried and appended twice. The fix restored acks=all and set enable.idempotence=true explicitly so the mistake cannot recur silently.
Second, even with idempotence on, a few duplicates remained. The exporter read pending invoices from a database table, sent them and then marked them as exported. Pods killed during the restart had sent some invoices but not yet marked them, so the replacement pods sent them again with a new producer ID. Idempotence cannot detect that; to the broker they are new records from a new producer. The real fix was on the consumer side: the ledger deduplicates by invoice ID, and the exporter puts that ID in a header and as the key.
ProducerRecord<String, String> rec = new ProducerRecord<>("invoices", invoice.id(), invoice.json());
rec.headers().add("event-id", invoice.id().getBytes(StandardCharsets.UTF_8));
producer.send(rec, callback);Common mistake
- “Idempotence gives exactly-once delivery.” It dedupes retries inside one producer session on one partition. Restarts, replays and two instances sending the same event still duplicate.
- “Idempotence deduplicates by key or content.” It uses PID and sequence numbers only; two identical records sent by the application are two records.
- “
send()returns after the broker acknowledges.” It returns aFutureimmediately after buffering. - “
linger.ms=0is lowest latency, so always use it.” Under load, tiny batches mean more requests and queueing, which can increase latency as well as cost. - “Set
max.in.flight.requests.per.connection=1for ordering.” Idempotence already preserves order with up to 5 in flight.
Verify the behavior
This toy broker in Node.js applies the same rule as the leader, deduplicating by producer ID, partition and sequence. Real brokers track per-record sequences per batch and remember only the last five batches; the model keeps just the last sequence. It printed the output in the comments when run with Node 22:
const log = [], lastSeq = new Map();
function append(pid, partition, seq, value) {
const k = `${pid}:${partition}`, last = lastSeq.get(k) ?? -1;
if (seq <= last) return 'DUPLICATE (acked, not appended)';
if (seq !== last + 1) return 'OUT_OF_ORDER_SEQUENCE (rejected)';
lastSeq.set(k, seq); log.push(value);
return `appended at offset ${log.length - 1}`;
}
console.log(append(1000, 0, 0, 'order-1:CREATED')); // appended at offset 0
console.log(append(1000, 0, 1, 'order-1:PAID')); // appended at offset 1 (ack lost)
console.log(append(1000, 0, 1, 'order-1:PAID')); // DUPLICATE (acked, not appended)
console.log(append(1000, 0, 3, 'order-1:REFUNDED')); // OUT_OF_ORDER_SEQUENCE (rejected)
console.log(append(1001, 0, 0, 'order-1:PAID')); // appended at offset 2 (new producer ID)The last line is the restart case: a new PID starts a fresh sequence, so the same event is stored twice. To measure batching on a real cluster, compare runs of the perf test tool:
kafka-producer-perf-test.sh --topic perf --num-records 1000000 --record-size 1000 \
--throughput -1 --producer-props bootstrap.servers=localhost:9092 linger.ms=0 batch.size=16384
kafka-producer-perf-test.sh --topic perf --num-records 1000000 --record-size 1000 \
--throughput -1 --producer-props bootstrap.servers=localhost:9092 linger.ms=20 batch.size=65536 compression.type=lz4Each run prints records per second, MB per second and latency percentiles. In the application, watch the producer metrics batch-size-avg, records-per-request-avg, record-queue-time-avg, compression-rate-avg and record-retry-rate.
Follow-up questions
What happens to idempotence if you set acks=1? If you did not set enable.idempotence, it is disabled silently; if you set it to true, the producer fails to start.
Why does send() sometimes block? Either buffer.memory is full because brokers are slower than the producer, or metadata for the topic is not available yet. Both wait up to max.block.ms.
How do you get deduplication across restarts? Use a transactional producer with a stable transactional.id, which fences the old instance, or deduplicate downstream by a business ID.
Which compression codec would you choose? lz4 for low CPU and good speed, zstd for the best ratio at moderate CPU, snappy for compatibility; gzip is the most CPU-expensive.
Interview exercise
A producer uses acks=all, idempotence on and linger.ms=0. Partition 0’s leader writes batch sequence 7, but the response is lost; meanwhile batch 8 is sent and acknowledged. The producer retries batch 7. Then the application crashes, restarts, and resends batch 8’s records. What does partition 0 contain, and how would you change the setup to raise throughput substantially?
Answer and reasoning
The retry of batch 7 is recognised as a duplicate: the broker already has sequence 7 for this producer ID and partition, so it acknowledges without appending. Batches 7 and 8 each appear once, in order. After the restart the producer has a new PID, so batch 8’s records are accepted again and appear twice; only consumer-side deduplication or transactions with a stable transactional.id prevent that. For throughput, raise linger.ms to around 10 to 20 ms, increase batch.size to 64 KB or more, enable lz4 or zstd compression, and make sure the code uses callbacks rather than .get() per record. Then verify with kafka-producer-perf-test.sh and the batch-size-avg metric instead of assuming the gain.
Continue learning
Practise with the Apache Kafka interview questions and the Kafka MCQs. Related notes: Kafka exactly-once and transactions, Kafka replication, ISR and acks and idempotency keys in microservices. Primary sources: the KafkaProducer Javadoc for Kafka 4.0, Confluent’s batch processing for efficiency and the producer configuration reference.