Menu

Kafka course · Lesson 6 of 6

Kafka Producers: acks, Batching, Compression, Idempotence and Ordering

How the Kafka producer sends records: acks and durability, idempotence, batching with linger.ms, compression codecs, partitioners, retries, in-flight requests and buffering.

  • Intermediate
  • 17 min read
  • Updated Oct 2026
On this page
  1. Sample setup
  2. Producer API
  3. The send path
  4. Worked example: Java
  5. Worked example: Python (confluent-kafka)
  6. Pitfalls
  7. In interviews
  8. acks: 0, 1 and all
  9. How it works with min.insync.replicas
  10. acks=0 on the cluster
  11. Pitfalls
  12. In interviews
  13. Idempotent producer
  14. Why duplicates happen without it
  15. How it works
  16. Worked example: a retry storm with and without idempotence
  17. Conflicting settings fail fast
  18. Pitfalls
  19. In interviews
  20. Batching and linger.ms
  21. How it works
  22. Worked example: measuring batching
  23. Tuning guidance
  24. Pitfalls
  25. In interviews
  26. Compression: snappy, lz4 and zstd
  27. How it works
  28. Worked example: same data, five codecs
  29. Pitfalls
  30. In interviews
  31. Partitioner strategies
  32. Pitfalls
  33. In interviews
  34. Retries and delivery
  35. How it works
  36. Pitfalls
  37. In interviews
  38. max.in.flight and ordering
  39. How it interacts with ordering
  40. Pitfalls
  41. In interviews
  42. Producer buffering
  43. How it works
  44. Pitfalls
  45. In interviews
  46. Practice questions
  47. Key takeaways

The producer looks simple: call send() and a record appears in a topic. Underneath, it batches records per partition, compresses them, sends several requests in parallel, retries on failure and waits for the acknowledgements you asked for. Each of those steps has a setting, and the defaults changed in Kafka 3.0 and 4.0. This lesson explains every one, so you can say exactly what your pipeline guarantees when a broker fails.

Sample setup

The examples write to the six-partition, replication-factor-3 orders topic from the topics lesson on a local three-broker Kafka 4.3.1 cluster:

bin/kafka-topics.sh --bootstrap-server localhost:19092 --create --topic orders \
  --partitions 6 --replication-factor 3

Java code needs the kafka-clients library (it is in the libs/ folder of the Kafka download, or from Maven Central). Python code uses the confluent-kafka package. The simulation further down needs only Python 3.

Producer API

What it is. The producer is a thread-safe client object. You configure it once, call send() (Java) or produce() (Python) for each record, and close it on shutdown. One producer per application process is usually right; it multiplexes every topic and partition over shared connections.

The send path

  1. Serialise the key and value to bytes with the configured serialisers.
  2. Partition: pick a partition (explicit, by key hash, or sticky for keyless records).
  3. Accumulate: append the record to an in-memory batch for that partition.
  4. A background sender thread drains ready batches, groups them by leader broker into produce requests, and sends them.
  5. The broker appends the batch and replies once the acks condition is met; the producer completes the record’s future and runs its callback.

send() is asynchronous: it returns as soon as the record is in a batch. You learn the outcome from the returned Future or the callback.

Worked example: Java

import org.apache.kafka.clients.producer.*;
import org.apache.kafka.common.serialization.StringSerializer;
import java.util.Properties;

public class OrderProducer {
    public static void main(String[] args) throws Exception {
        Properties props = new Properties();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:19092");
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        props.put(ProducerConfig.CLIENT_ID_CONFIG, "order-service");
        // Defaults since Kafka 3.0, set explicitly so intent is visible in code review
        props.put(ProducerConfig.ACKS_CONFIG, "all");
        props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "true");
        props.put(ProducerConfig.LINGER_MS_CONFIG, "20");
        props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, "zstd");

        try (Producer<String, String> producer = new KafkaProducer<>(props)) {
            String[][] events = {
                {"cust-1", "{\"order_id\":1001,\"status\":\"created\"}"},
                {"cust-2", "{\"order_id\":1002,\"status\":\"created\"}"},
                {"cust-1", "{\"order_id\":1001,\"status\":\"paid\"}"},
            };
            for (String[] e : events) {
                ProducerRecord<String, String> record = new ProducerRecord<>("orders", e[0], e[1]);
                producer.send(record, (metadata, exception) -> {
                    if (exception != null) {
                        System.err.println("send failed for key " + e[0] + ": " + exception);
                    } else {
                        System.out.printf("key=%s -> partition %d, offset %d%n",
                                e[0], metadata.partition(), metadata.offset());
                    }
                });
            }
            producer.flush();
        }
    }
}
key=cust-1 -> partition 1, offset 6
key=cust-2 -> partition 1, offset 7
key=cust-1 -> partition 1, offset 8

Worked example: Python (confluent-kafka)

The Confluent client wraps librdkafka. Delivery reports only fire when you call poll() or flush().

import json
from confluent_kafka import Producer

producer = Producer({
    "bootstrap.servers": "localhost:19092",
    "client.id": "order-service-py",
    "acks": "all",
    "enable.idempotence": True,
    "linger.ms": 20,
    "compression.type": "zstd",
    "partitioner": "murmur2_random",   # match the Java client's key hashing
})

def on_delivery(err, msg):
    if err is not None:
        print(f"FAILED key={msg.key()}: {err}")
    else:
        print(f"key={msg.key().decode()} -> partition {msg.partition()}, offset {msg.offset()}")

events = [("cust-1", {"order_id": 1001, "status": "refunded"}),
          ("cust-4", {"order_id": 1004, "status": "created"})]
for key, value in events:
    producer.produce("orders", key=key, value=json.dumps(value), on_delivery=on_delivery)
    producer.poll(0)          # serve delivery callbacks for earlier sends

remaining = producer.flush(10)  # block until everything is acknowledged or 10 s pass
print("undelivered:", remaining)
key=cust-4 -> partition 0, offset 1620
key=cust-1 -> partition 1, offset 2710
undelivered: 0

(The offsets are high because earlier experiments had written to the topic.)

Pitfalls

  • Fire and forget by accident. Calling send() without a callback or checking the future means failures are only logged. Always handle the result, and count failures in a metric.
  • Calling .get() on every send makes each send synchronous, so batching cannot happen. Use callbacks, or .get() only at a checkpoint.
  • Not flushing or closing on shutdown. Records still in the accumulator are lost when the process exits. close() flushes; in Python call flush().
  • One producer per record or per request. Creating producers is expensive (connections, metadata, producer ID). Reuse one.

In interviews

Walk through serialise, partition, accumulate, send, acknowledge, and say that send() is asynchronous. Mention that the producer is thread-safe and should be shared.

acks: 0, 1 and all

What it is. acks says how many replicas must have the record before the leader replies “success”.

acks Leader replies when Loss risk Latency
0 never (producer does not wait) Highest: any failure loses data silently Lowest
1 the leader has written it to its log Leader dies before followers copy it: acknowledged data lost Low
all (-1) every replica in the current in-sync set (ISR) has it Only if all ISR replicas are lost Highest

acks=all has been the default since Kafka 3.0 (it was 1 before).

How it works with min.insync.replicas

acks=all waits for the current ISR. If the ISR has shrunk to just the leader, acks=all is no safer than acks=1. The topic or broker setting min.insync.replicas closes that gap: when the ISR is smaller than it, the leader rejects acks=all writes with NotEnoughReplicasException instead of accepting them with too few copies. The standard durable combination is replication factor 3, min.insync.replicas=2, acks=all: one broker can be down and writes continue with two copies. The broker default min.insync.replicas is 1, so set it deliberately. The replication lesson shows this failing on a real cluster.

acks=0 on the cluster

With acks=0 the producer cannot know the offset, so the metadata it returns has offset -1:

acks=0: ok, offset -1

Pitfalls

  • acks=all with min.insync.replicas=1 gives a false sense of safety.
  • acks=0 also disables retries in effect: the producer never sees errors.
  • Requiring min.insync.replicas equal to the replication factor means any single broker restart stops writes.

In interviews

“What does acks=all guarantee?” Strong answer: the record is on every in-sync replica before success is returned, which only means something if min.insync.replicas is at least 2. Mention the 3.0 default change.

Idempotent producer

What it is. An idempotent producer can retry without creating duplicates or reordering records within a partition. It is enabled by default since Kafka 3.0 (when no conflicting settings are given).

Why duplicates happen without it

A produce request can succeed on the broker while the response is lost (network blip, timeout). The producer cannot tell this from a failure, so it retries, and the record is written twice. Retries can also reorder: if batch 1 fails and batch 2 (already in flight) succeeds, the retry of batch 1 lands after batch 2.

How it works

  • On start-up the producer gets a producer ID (PID) and epoch from the cluster.
  • Every batch carries the PID and a per-partition sequence number. You saw these in the kafka-dump-log.sh output in the log storage lesson (producerId: 1 ... baseSequence: 1984).
  • The partition leader remembers the last sequence numbers for each PID. A batch with a sequence it has already written is acknowledged again without being appended (a duplicate). A batch that skips ahead is rejected as out of order, and the producer resends in sequence.
  • Requirements: acks=all, retries > 0, and max.in.flight.requests.per.connection <= 5, because the broker keeps state for at most five batches per producer per partition.

Worked example: a retry storm with and without idempotence

The model below plays the same failure against a broker that ignores sequence numbers and one that enforces them. Batch B0 is written but its acknowledgement is lost, B1 fails with a transient error, and B2 arrives while both are unresolved.

class Partition:
    """A partition leader that optionally enforces idempotent-producer sequence numbers."""
    def __init__(self, idempotent):
        self.idempotent = idempotent
        self.log = []
        self.next_seq = {}          # producer id -> next expected sequence

    def append(self, pid, seq, batch):
        if self.idempotent:
            expected = self.next_seq.get(pid, 0)
            if seq < expected:
                return "DUPLICATE (already written, ack again)"
            if seq > expected:
                return "OUT_OF_ORDER_SEQUENCE (rejected, retry later)"
            self.next_seq[pid] = seq + 1
        self.log.append(batch)
        return "OK"

def run(idempotent):
    leader = Partition(idempotent)
    pid = 42
    pending = [(0, "B0"), (1, "B1"), (2, "B2")]   # (sequence, batch), up to 5 in flight
    # Round 1: all three are in flight. B0's write succeeds but its ack is lost on the network;
    # B1 hits a transient error and is not written; B2 arrives and is processed.
    outcomes = {}
    leader.append(pid, 0, "B0"); outcomes["B0"] = "written, ack lost -> retry"
    outcomes["B1"] = "transient error -> retry"
    outcomes["B2"] = leader.append(pid, 2, "B2")
    # Round 2: the producer retries every batch that was not acknowledged, in sequence order.
    retries = {}
    for seq, batch in pending:
        if batch == "B2" and outcomes["B2"] == "OK":
            continue
        retries[batch] = leader.append(pid, seq, batch)
    print("idempotent" if idempotent else "not idempotent")
    print("  first attempt:", outcomes)
    print("  retries:      ", retries)
    print("  partition log:", leader.log)

run(idempotent=False)
run(idempotent=True)
not idempotent
  first attempt: {'B0': 'written, ack lost -> retry', 'B1': 'transient error -> retry', 'B2': 'OK'}
  retries:       {'B0': 'OK', 'B1': 'OK'}
  partition log: ['B0', 'B2', 'B0', 'B1']
idempotent
  first attempt: {'B0': 'written, ack lost -> retry', 'B1': 'transient error -> retry', 'B2': 'OUT_OF_ORDER_SEQUENCE (rejected, retry later)'}
  retries:       {'B0': 'DUPLICATE (already written, ack again)', 'B1': 'OK', 'B2': 'OK'}
  partition log: ['B0', 'B1', 'B2']

Without idempotence the log has a duplicate (B0 twice) and B1 after B2. With it, the log is exactly B0, B1, B2.

Conflicting settings fail fast

If you explicitly enable idempotence and set something incompatible, Kafka 4.x refuses to start the producer (older versions sometimes silently disabled idempotence):

acks=1 + enable.idempotence=true: ConfigException: Must set acks to all in order to use the idempotent producer. Otherwise we cannot guarantee idempotence.
max.in.flight=6 + enable.idempotence=true: ConfigException: To use the idempotent producer, max.in.flight.requests.per.connection must be set to at most 5. Current value is 6.

If you do not set enable.idempotence but set a conflicting value such as acks=1, idempotence is quietly turned off instead.

Pitfalls

  • Scope. Idempotence covers one producer session writing to Kafka. If your application restarts and resends the same business event, that is a new write with a new PID. End-to-end deduplication needs transactions or an idempotent sink (see delivery semantics).
  • Setting acks=1 for speed silently removes idempotence unless you also set enable.idempotence=true, in which case the producer will not start.

In interviews

Explain PID plus per-partition sequence numbers, the broker’s duplicate and out-of-order checks, the three requirements, and that it is on by default since 3.0. Then state its limit: it does not survive an application-level resend.

Batching and linger.ms

What it is. The producer groups records for the same partition into a batch and sends batches, not records. Bigger batches mean fewer requests, better compression and higher throughput, at the cost of a little latency.

How it works

  • batch.size (default 16384 bytes) is the target maximum batch size per partition. A full batch is sent immediately.
  • linger.ms is how long the producer waits for more records before sending a batch that is not full. The default changed from 0 to 5 ms in Kafka 4.0.
  • Under load, batching happens even with linger.ms=0, because records pile up while earlier requests are in flight.
  • One produce request carries batches for several partitions led by the same broker.

Worked example: measuring batching

This test sent 1,000 small records at roughly one per millisecond to a single-partition topic with replication factor 3, then read the producer’s own metrics (batch-size-avg, records-per-request-avg, request-total):

linger.ms=0   batch-size-avg=  314.9 bytes  records-per-request-avg=  11.0  request-total=96
linger.ms=10  batch-size-avg=  316.4 bytes  records-per-request-avg=  11.2  request-total=93
linger.ms=100 batch-size-avg= 1278.4 bytes  records-per-request-avg=  52.6  request-total=24

A second run gave different absolute numbers (for example 146 requests at linger.ms=0 and 18 at 100), but the same shape. Even linger.ms=0 batched about ten records per request, because each acks=all round trip took a few milliseconds while new records arrived. Raising linger.ms to 100 cut requests by roughly four times. On your own system, read these metrics rather than guessing.

Tuning guidance

Goal Settings to try
High-throughput ingestion linger.ms 10 to 100, batch.size 64 KiB to 256 KiB, compression on
Low latency per event keep linger.ms small (the default 5 is fine), compression lz4 or none
Many partitions, low rate per partition expect small batches; consider fewer partitions or a larger linger.ms

Pitfalls

  • A huge batch.size with many partitions reserves a lot of buffer memory, because each partition can have an open batch.
  • Keyed traffic spread over many partitions produces small batches per partition; batching works per partition.
  • Confusing linger.ms with end-to-end latency. It is an upper bound on the extra wait before sending, not a delay applied to every record when batches fill quickly.

In interviews

“How do you increase producer throughput?” Mention linger.ms and batch.size, compression, async sends with callbacks, and checking batch-size-avg and records-per-request-avg. Note the 4.0 default of 5 ms.

Compression: snappy, lz4 and zstd

What it is. The producer can compress each batch with gzip, snappy, lz4 or zstd (compression.type, default none). The broker stores the batch compressed and consumers decompress it, so compression saves network, disk and replication traffic.

How it works

  • Compression is per batch, so bigger batches compress better.
  • The topic setting compression.type defaults to producer, meaning “keep whatever codec the producer used”. If you set a topic to a specific codec and it differs from the producer’s, the broker recompresses, which costs broker CPU.
  • Levels can be tuned with compression.gzip.level, compression.lz4.level and compression.zstd.level (defaults -1, 9 and 3 in Kafka 4.3).

Worked example: same data, five codecs

20,000 JSON click events (about 100 bytes each, generated from a fixed random seed) were written to five single-partition topics with linger.ms=50 and batch.size=65536, one codec per topic. On-disk partition size reported by kafka-log-dirs.sh:

compression.type Bytes on disk Relative to none
none 2,039,107 100%
snappy 532,842 26%
lz4 523,289 26%
gzip 322,660 16%
zstd 308,926 15%

That is one small, repetitive dataset; your ratios depend on your data and batch sizes, and this test did not measure CPU. The usual trade-off: lz4 and snappy are very fast with moderate ratios; zstd gives gzip-like ratios at much lower CPU cost than gzip; gzip is the slowest.

Pitfalls

  • Compressing tiny batches gains little. Fix batching first.
  • Consumers pay to decompress. Very old clients may not support zstd (it arrived in Kafka 2.1).
  • Already-compressed payloads (images, Parquet, Avro with its own codec) do not shrink further.

In interviews

Say compression is per batch, is applied by the producer, and the broker keeps it as is when the topic uses compression.type=producer. A sensible default answer is lz4 or zstd.

Partitioner strategies

What it is. The partitioner chooses a partition for each record that does not name one.

Strategy When used Behaviour
Explicit partition new ProducerRecord<>(topic, partition, key, value) Uses the given partition
Key hash (default with key) key is not null murmur2 hash of the key bytes modulo partition count
Sticky (default without key) key is null Fills a batch for one partition, then switches; since Kafka 3.3 it also favours faster brokers (partitioner.adaptive.partitioning.enable=true)
partitioner.ignore.keys=true you want spread, not key ordering Treat keyed records like keyless ones
RoundRobinPartitioner partitioner.class Each record to the next partition, ignoring keys
Custom Partitioner partitioner.class Your own logic, for example routing a VIP tenant to dedicated partitions

The old DefaultPartitioner and UniformStickyPartitioner classes were deprecated in 3.3 and are no longer in the Kafka 4.x client library; leave partitioner.class unset to get the built-in behaviour.

Pitfalls

  • Clients disagree on hashing. librdkafka-based clients default to a CRC32 partitioner; set partitioner=murmur2_random to match Java (see the topics lesson for a test).
  • Custom partitioners that depend on partition count break key ordering when partitions are added, just like the default.
  • Hot keys. No partitioner can split one key across partitions without giving up its ordering. If one key dominates, add a sub-key (for example customer_id plus a bucket) only if per-customer order does not matter.

In interviews

Know the three defaults (explicit, key hash, sticky) and when you would write a custom partitioner. Mention that sticky partitioning replaced round-robin for keyless data because it produces larger batches.

Retries and delivery

What it is. Transient errors (leader moved, not enough replicas, request timeout) are retried automatically. The modern way to bound retries is by time, not count.

How it works

Setting Default Meaning
retries 2147483647 effectively unlimited; leave it
delivery.timeout.ms 120000 (2 min) total time from send() returning to success or failure, including batching, waiting and all retries
request.timeout.ms 30000 how long to wait for one response before retrying
retry.backoff.ms / retry.backoff.max.ms 100 ms / 1000 ms exponential back-off between attempts

delivery.timeout.ms must be at least linger.ms + request.timeout.ms. When it expires, the callback receives a TimeoutException and the record is not retried further: your code must decide what to do (log, send to a dead-letter store, fail the job).

Non-retriable errors fail at once, for example a record larger than max.request.size (1 MiB by default):

RecordTooLargeException: The message is 1100089 bytes when serialized which is larger than 1048576, which is the value of the max.request.size configuration.

Pitfalls

  • Lowering retries to avoid duplicates. That trades duplicates for data loss; keep idempotence on instead.
  • Swallowing callback errors. After delivery.timeout.ms the data is gone unless you handle the exception.
  • Large records need the producer’s max.request.size, the topic’s max.message.bytes and consumer fetch limits raised together, or better, store the payload elsewhere and send a reference.

In interviews

Explain that retries are bounded by delivery.timeout.ms, that retries without idempotence can duplicate and reorder, and what your code does with a final failure.

max.in.flight and ordering

What it is. max.in.flight.requests.per.connection (default 5) is how many unacknowledged requests the producer may have outstanding to one broker. More in flight means higher throughput on high-latency links.

How it interacts with ordering

Idempotence max.in.flight Ordering on retry
on (default) 1 to 5 Preserved: the broker rejects out-of-sequence batches
off 1 Preserved, but throughput is limited
off more than 1 Can reorder: a retried batch lands after later ones
on more than 5 Not allowed: ConfigException in Kafka 4.x

Before idempotence existed, the standard advice for strict ordering was max.in.flight.requests.per.connection=1. With the idempotent producer, you keep up to 5 and still get ordering.

Pitfalls

  • Copying old tuning guides that set max.in.flight=1 “for ordering” costs throughput for no benefit when idempotence is on.
  • Raising it above 5 with idempotence enabled now fails; without explicit idempotence it silently disables idempotence.

In interviews

“How do you keep order with retries?” Say: idempotent producer on (the default), max.in.flight at most 5, keyed records, and one partition per key. Mention max.in.flight=1 as the pre-idempotence answer.

Producer buffering

What it is. Records wait in the record accumulator, an in-memory buffer of size buffer.memory (default 32 MiB), until the sender thread ships them. If the application produces faster than brokers accept, the buffer fills.

How it works

  • When the buffer is full, send() blocks for up to max.block.ms (default 60 s) waiting for space. The same limit covers waiting for topic metadata.
  • If space does not appear in time, send() fails with a TimeoutException; the error surfaces through the returned future or callback.
  • A single record larger than buffer.memory or max.request.size fails immediately.
  • In the Python client the equivalent queue limits are queue.buffering.max.messages and queue.buffering.max.kbytes; a full queue raises BufferError from produce(), and the fix is to call poll() and retry.

Pitfalls

  • Blocking a request thread. In a web service, send() blocking for 60 seconds when Kafka is unavailable can exhaust the server’s threads. Lower max.block.ms and decide whether to drop, spill to disk or fail fast.
  • Unbounded memory in your own code. If you queue records before calling send(), you have built a second, unmonitored buffer.
  • Watch the metrics buffer-available-bytes and bufferpool-wait-time-ns-total (the Java producer’s metrics) to see back-pressure before it becomes an outage.

In interviews

“What happens when the broker is slow or down?” Records accumulate, send() blocks up to max.block.ms once buffer.memory is full, unsent records expire after delivery.timeout.ms, and your callback must handle both failures.

Practice questions

Your producer uses acks=all but you still lost acknowledged data after a broker failure. How?

acks=all waits only for the current in-sync replicas. If min.insync.replicas was 1 and the ISR had shrunk to the leader alone, the leader acknowledged writes on its own; when it failed, an out-of-date replica could not have them (and with unclean leader election enabled, it might even become leader). Use replication factor 3 with min.insync.replicas=2 and keep unclean leader election disabled.

What problem does the idempotent producer solve, and what can it not solve?

It stops retries from writing duplicates or reordering batches within a partition, using a producer ID and per-partition sequence numbers checked by the broker. It cannot deduplicate an application that crashes and resends the same business event with a new producer session, and it says nothing about consumers processing a record twice. Those need transactions or idempotent consumers and sinks.

A team sets max.in.flight.requests.per.connection=1 to keep ordering. Is that still needed?

No, if idempotence is enabled (the default since 3.0). The broker rejects out-of-sequence batches, so ordering is preserved with up to five requests in flight. Setting 1 only reduces throughput. It is still needed for ordering if idempotence is disabled.

How would you raise producer throughput for a bulk backfill?

Send asynchronously with callbacks, raise linger.ms (for example 50 to 100 ms) and batch.size (for example 128 KiB), enable lz4 or zstd compression, make sure buffer.memory is large enough for the open batches, and check batch-size-avg, records-per-request-avg and record-send-rate. If keys spread data over many partitions with tiny batches, consider whether the topic needs that many partitions.

What happens to a record when the brokers are unreachable for five minutes?

send() may block for up to max.block.ms waiting for metadata or buffer space. Records already accepted are retried until delivery.timeout.ms (two minutes by default) expires, then their callbacks receive a TimeoutException. Unless the application stores or resends them, they are lost.

Why does the broker sometimes spend CPU recompressing data?

When the topic’s compression.type is set to a specific codec that differs from the producer’s codec, the broker must decompress and recompress each batch. Leaving the topic at compression.type=producer lets the broker store batches as they arrive.

Key takeaways

  • send() is asynchronous: serialise, partition, batch, send, acknowledge. Always handle the callback or future.
  • Use acks=all (the default since 3.0) with min.insync.replicas=2 on replication-factor-3 topics for durability.
  • The idempotent producer, on by default since 3.0, removes retry duplicates and reordering within a partition, but only within one producer session.
  • linger.ms (5 ms by default since 4.0), batch.size and compression drive throughput; measure with producer metrics.
  • Retries are bounded by delivery.timeout.ms; final failures reach your callback and need a decision.
  • When brokers fall behind, buffer.memory fills and send() blocks for max.block.ms; plan for back-pressure.

By Data Career Hub Editorial · Last reviewed Oct 2026 · Java producer examples were compiled and run against a three-broker Apache Kafka 4.3.1 KRaft cluster with kafka-clients 4.3.1 on Java 21; the Python producer was run with confluent-kafka 2.15.1 (librdkafka 2.15.1). They are marked noexec because they need a running broker. Defaults come from the Kafka 4.3.1 configuration reference. The idempotence simulation runs on plain Python 3.

Progress is saved in this browser only. No account needed.

Search
Filter by type