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

Quick Example

A durable, efficient Java producer:

Core Concepts

The Send Path

  1. Serialize key and value to bytes.
  2. Partition: choose a partition from the key's hash, or batch keyless records to a partition.
  3. Accumulate: append to an in-memory batch for that partition (bounded by buffer.memory; send() blocks up to max.block.ms if full).
  4. Send: a background I/O thread sends batches to partition leaders, up to max.in.flight.requests.per.connection concurrently per broker.
  5. 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

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:

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

References