Kafka Producers
A producer publishes records to Kafka topics. On the surface it's a one-line send() call, but underneath, the client buffers records, groups them into batches per partition, compresses them, sends them to the partition leaders, and retries on failure. How you configure that pipeline decides whether messages can be lost, duplicated, or reordered, and how much throughput you get.
Modern Kafka clients default to safe settings: acks=all and idempotence enabled. Understanding what those settings mean, and how batching and compression trade latency for throughput, lets you tune producers for your workload without sacrificing correctness.
TL;DR
send()is asynchronous: records are buffered, batched per partition, and sent in the background. Handle the result in a callback or future.acks=allplus the topic'smin.insync.replicasmakes writes durable across broker failures.- Idempotent producers (
enable.idempotence=true, the default) prevent duplicates and reordering caused by retries. - Retries are bounded by
delivery.timeout.ms. Treat a final send failure as a real error. linger.ms,batch.size, andcompression.typetrade a few milliseconds of latency for much higher throughput.- Use a key for ordering, and a schema (Avro, Protobuf, JSON Schema) with a registry for safe evolution.
Quick Example
A durable, efficient Java producer:
Core Concepts
The Send Path
- Serialize key and value to bytes.
- Partition: choose a partition from the key's hash, or batch keyless records to a partition.
- Accumulate: append to an in-memory batch for that partition (bounded by
buffer.memory;send()blocks up tomax.block.msif full). - Send: a background I/O thread sends batches to partition leaders, up to
max.in.flight.requests.per.connectionconcurrently per broker. - Acknowledge: the broker responds according to
acks, and the callback or future completes with metadata or an exception.
acks and Durability
acks=all alone isn't enough: if the ISR shrinks to just the leader, "all" means one copy. Set min.insync.replicas=2 on the topic, so writes fail rather than proceed with a single copy. See Kafka topics & partitions.
Idempotence
Without idempotence, a retry after a lost acknowledgement writes the record twice, and retries with multiple in-flight requests can reorder records. An idempotent producer gets a producer ID and attaches sequence numbers per partition; brokers discard duplicates and reject out-of-order sequences. It's enabled by default in modern clients (Kafka 3.0+) when acks=all, and it keeps ordering with up to 5 in-flight requests. It guarantees exactly-once writes per partition per producer session. For atomic writes across partitions, use transactions.
Retries and Timeouts
Transient errors (leader elections, network blips, NOT_ENOUGH_REPLICAS) are retried automatically. The total time a record may spend being retried is bounded by delivery.timeout.ms (default 2 minutes). After that, the send fails with an exception in the callback. That final failure must be handled: log and alert, write to a fallback store, or propagate to the caller.
Batching and Compression
linger.ms: how long to wait for more records before sending a batch (default 5 ms in recent clients). Higher values mean bigger batches, better throughput and compression, and slightly more latency.batch.size: the maximum bytes per partition batch.compression.type:zstdorlz4usually give the best ratio and speed trade-off, andsnappyandgzipare available. Compression works per batch, so larger batches compress better.
Throughput tuning is mostly about making batches bigger: raise linger.ms and batch.size, enable compression, and send from fewer, longer-lived producer instances.
Serialization and Schemas
Producers turn objects into bytes with serializers. For long-lived topics shared across teams, use a schema format (Avro, Protobuf, or JSON Schema) with a Schema Registry: producers register schemas, records carry a schema ID, and compatibility rules (for example backward-compatible) prevent a producer from breaking consumers. Plain JSON without schemas works for small systems but makes evolution risky.
Reliable Publishing Patterns
A common failure mode: the service commits a database transaction, then crashes before publishing the event (or publishes, then the transaction rolls back). Solutions:
- Transactional outbox: write the event to an
outboxtable in the same database transaction, and let a relay or CDC connector (Debezium) publish it to Kafka. - Kafka transactions when the source of truth is Kafka itself (consume-transform-produce).
- Idempotent consumers downstream, because at-least-once publishing still happens in some failure paths. See idempotency.
Best Practices
Reuse One Producer per Application
Producers are thread-safe and expensive to create (connections, buffers, metadata). Create one per process, or per distinct configuration, and share it. Close it on shutdown with flush() so buffered records aren't lost.
Never Ignore Send Results
Fire-and-forget send(record) without a callback silently drops failures after retries are exhausted. Always check the future or callback, and export error-rate metrics.
Keep Keys Meaningful and Stable
The key determines partition and ordering. Use the business entity whose events must stay ordered, and keep the key format consistent across producers. "42" and 42 serialize differently and land in different partitions.
Monitor Producer Metrics
Watch record-error-rate, record-retry-rate, request-latency-avg, batch-size-avg, compression-rate-avg, and buffer-available-bytes. A shrinking buffer means the producer can't keep up with brokers.
Common Mistakes
acks=1 for Important Data
Use acks=all with min.insync.replicas=2 for anything you can't afford to lose.
Creating a Producer per Request
Instantiating a KafkaProducer inside a request handler adds connection setup to every call, defeats batching, and can exhaust broker connections. Share one producer.
Blocking on Every Send
Calling producer.send(record).get() for each message makes throughput latency-bound, one round trip per record. Send asynchronously and handle results in callbacks. When you need confirmation, wait on a batch of futures.
FAQ
Does Kafka guarantee no duplicates?
With idempotence enabled (the default), retries within a producer session won't create duplicates in a partition. Duplicates can still arise at the application level, for example when your service restarts and re-sends a record it didn't know was written. Use transactions or idempotent consumers for end-to-end exactly-once semantics.
How do I maximize producer throughput?
Increase batching (linger.ms 10–50 ms, larger batch.size), enable zstd or lz4 compression, send asynchronously, reuse producers, and make sure topics have enough partitions spread across brokers. Measure with kafka-producer-perf-test.sh.
Does message order hold with retries?
Yes, with idempotence enabled: sequence numbers let brokers reject out-of-order batches, so per-partition order is preserved with up to 5 in-flight requests. Without idempotence, retries with max.in.flight.requests.per.connection > 1 can reorder records.
Should I use JSON or Avro/Protobuf?
For shared, long-lived topics, a schema format with a registry is strongly recommended: compact binary encoding, enforced compatibility, and generated types. JSON is fine for prototypes and small, single-team systems, ideally still validated against a JSON Schema.
Related Topics
- Kafka — The platform overview
- Kafka Topics & Partitions — Keys, replication, and min.insync.replicas
- Kafka Consumer Groups — The other side of the pipe
- Kafka Exactly-Once — Transactions and atomic writes
- Change Data Capture — The outbox pattern via CDC
- Idempotency — Handling duplicates safely