Skip to content

Kafka Producer Tuning Cheat Sheet: Throughput, Latency & Durability

Kafka Producer Tuning Cheat Sheet: Throughput, Latency & Durability

A broker fails over at peak, and 30 seconds of in-flight messages vanish — no error, no retry. The config: acks=1, linger.ms=0, retries=0, enable.idempotence=false — it passes staging fine. But acks=1 only confirms the leader, not the followers, so a leader failover silently drops everything in flight.

That config predates Kafka 3.0, which made the Java producer safe-by-default (enable.idempotence=true, acks=allKIP-679). You reach the failure mode today through legacy configs, explicit overrides, or non-Java clients with weaker defaults — verify every value below rather than trusting an inherited config.

Kafka Producer Tuning Essentials

TL;DR

Start with: acks=all, enable.idempotence=true, batch.size=65536 (up from the 16384 default), linger.ms=5 (default 0), compression.type=lz4. Idempotence costs low-single-digit throughput per Confluent's published benchmarks (source) — measure on your own record sizes when it matters. Document every deviation.

  • Survive broker failover: acks=all + broker min.insync.replicas=2 (3-replica cluster)
  • Biggest throughput lever: linger.ms — raise from 0 to 5–20ms, not batch.size
  • One KafkaProducer per app; reuse it thread-safe, never create per-message
graph LR
    App["Application<br/>producer.send()"] --> RA["RecordAccumulator<br/>buffer per partition"]
    RA -->|"batch.size or<br/>linger.ms expires"| ST["Sender Thread"]
    ST -->|"max.in.flight<br/>batches"| Broker["Kafka Broker"]
    Broker -->|"acks=all"| ST
    ST -->|"callback"| App

The Quick Start

Goalacksbatch.sizelinger.mscompressionmax.in.flightTrade-off
Durability first (default)all65536 (64K)5lz45Loses ~5ms latency, gains no data loss
Maximum throughputall131072 (128K)20zstd5Adds ~20ms latency, 15% compression gain
Latency-critical (<10ms p99)all163840none1Drops to single partition throughput, keeps durability
Metrics you can lose0655365lz45Fastest; data loss on any broker crash

Durability: Choosing acks

[Kafka producer config]

acks is the kill switch. acks=0 loses data on restart. acks=1 loses data on leader failover — it looks safe in tests and fails silently in production. Use acks=all everywhere except metrics.

Pair it with broker-side min.insync.replicas=2 on a 3-replica cluster. This is the only combination that survives a broker failure without data loss.

Throughput: Batching + Linger

[Apache Kafka Docs]

Three knobs move throughput: batch.size, linger.ms, and compression.type. batch.size is a ceiling per partition (default 16384 bytes). linger.ms is the biggest lever most teams miss — raising from 0 to 5–20ms can 4x throughput on small records. Use lz4 (good ratio/CPU) unless you're network-bound; zstd achieves better compression ratios but costs more CPU.

Measure batch fill ratio: record-size-avg / batch-size-avg. Below 0.5 means raise linger.ms.

Buffer + sender thread lifecycle

send() appends to an in-memory RecordAccumulator; a background Sender thread drains per-partition batches when they fill OR linger.ms expires.

sequenceDiagram
    participant App
    participant Acc as Accumulator
    participant Sender
    participant Broker

    App->>Acc: send(record)
    App->>Acc: send(record)
    App->>Acc: send(record)
    Note over Acc: batch fills OR<br/>linger.ms expires
    Acc->>Sender: flush batch
    Sender->>Broker: Produce (acks=all)
    Broker-->>Sender: ack (ISR copied)
    Sender-->>App: onCompletion

If buffer.memory fills before the Sender can drain, send() blocks until space frees. This is the hidden backpressure path — misconfigured linger.ms + slow broker can stall the application thread.

Idempotence & Ordering

[Kafka producer config]

Enable enable.idempotence=true unconditionally. It adds 5 bytes per record and eliminates duplicates on retry. With idempotence, max.in.flight.requests.per.connection stays at 5; without it, drop to 1 for ordering.

Order is guaranteed within a partition. Key records by entity ID for strict per-entity ordering:

producer.send(new ProducerRecord<>("orders", order.customerId(), order));

Don't use transactions for ordering — they're for atomic multi-partition writes and exactly-once semantics, and cut throughput 20%. [Kafka producer config]

Tune by symptom

[Kafka producer config]

When a Kafka producer is misbehaving, the question is "what is the symptom?" — not "which config knob shall I tweak?" Route by what you measured:

graph TD
    Sym{What is<br/>broken?} -->|Data loss on failover| Loss[acks=all<br/>+ min.insync.replicas=2<br/>+ enable.idempotence=true]
    Sym -->|Low throughput<br/>under 10k msg/s| Tp[linger.ms 5 to 20<br/>+ batch.size 65536<br/>+ compression.type=lz4]
    Sym -->|High p99 latency| Lat[linger.ms back to 0 to 1<br/>+ acks=1 if data-loss tolerable<br/>+ smaller batch.size]
    Sym -->|Duplicates downstream| Dup[enable.idempotence=true<br/>+ transactional.id for exactly-once]
    Sym -->|Out-of-order messages| Ord[max.in.flight.requests=1<br/>or enable.idempotence=true<br/>which keeps order with five in-flight]
    Sym -->|Producer blocks on send| Buf[buffer.memory 256 MiB or higher<br/>+ check broker backpressure]
    Sym -->|Network-blip outages| Net[retries=MAX_INT<br/>+ delivery.timeout.ms 120 seconds<br/>+ enable.idempotence=true]
    style Loss fill:#fdd
    style Dup fill:#ffd
    style Ord fill:#ffd
    style Buf fill:#fdd
    style Net fill:#dfd
    style Tp fill:#dfd
    style Lat fill:#dfd

Classify the symptom, pick the knobs that apply, never tune in isolation.[Kafka producer config]

Common Gotchas

  • acks=1 + min.insync.replicas=1: Loses data on failover silently. Use acks=all with min.insync.replicas=2.
  • retries=0: Drops on network blips. Use enable.idempotence=true + retries=Integer.MAX_VALUE.
  • 32 MiB buffer on high-throughput: Fills at 10k msg/s, blocks send(). Raise to 256 MiB+.
  • linger.ms=0 everywhere: Default is 0 for legacy reasons. Set to 5–20ms on non-latency-critical producers.
  • Per-message KafkaProducer: The client is thread-safe. Create one per app, never per-request; causes metadata fetch storms.

Producer Metrics Worth Alerting On

Five JMX/MBean metrics that make production producer issues visible before they become incidents[Apache Kafka Docs]:

MetricThresholdWhat it meansFix
record-error-rate> 0.001 (0.1%)Send failures (after retries) — broker rejecting or unrecoverableCheck broker logs; verify ACL; raise delivery.timeout.ms
record-queue-time-avg> linger.ms × 2Records waiting in producer buffer too longbatch.size too small OR broker under-acking; profile broker
record-send-rate vs record-error-rateerror / send > 0.005Cluster instability or topic mis-configCheck metadata-fetch-rate for storms; verify min.insync.replicas
request-latency-avg> 100 ms p99Network or broker slownesstcpdump between producer and broker; check broker GC pauses
buffer-available-bytes / buffer-total-bytes< 0.10 (10%)Buffer about to fill; send() will blockRaise buffer.memory; check downstream broker backpressure

Export via the JMX exporter (JVM) or the client's metrics() map (Go). Burn-rate alert on record-error-rate; saturation alert on buffer-available-bytes.

Sizing Worker Threads in Java Apps

The producer is non-blocking by default — send() returns a Future immediately and the I/O thread handles delivery. Application threads should NOT wait on the future synchronously; use a callback or accumulate batches:

// Anti-pattern: blocks application thread per message
producer.send(record).get();  // synchronous — defeats batching
 
// Production pattern: callback + structured error handling
producer.send(record, (metadata, exception) -> {
    if (exception != null) {
        log.error("send failed: topic={}, partition={}", record.topic(), record.partition(), exception);
        deadLetterQueue.offer(record);  // application-side fallback
    }
});

In Go (sarama/franz-go): async-producer with a goroutine draining Successes and Errors — never silently drop the error channel.

Production Tuning Checklist

Apply in this order — each step depends on the prior:

  1. Durability first: acks=all, min.insync.replicas=2 on a 3-replica topic. Confirm with kafka-configs.sh --describe.
  2. Idempotence: enable.idempotence=true. Free in Kafka 3.x — do this before any other tuning.
  3. Throughput second: linger.ms=10, batch.size=65536, compression.type=lz4. Measure record-send-rate before/after.
  4. Buffer + delivery timeouts: buffer.memory=268435456 (256 MiB), delivery.timeout.ms=120000. Prevents send() blocking under broker stalls.
  5. Per-tenant quotas (multi-tenant clusters): Configure broker-side quotas via kafka-configs.sh --add-config 'producer_byte_rate=...'. Producer-side enforcement is fragile.
  6. Observability: Wire the 5 metrics above to Prometheus + alerting. Verify with jconsole or an equivalent for non-JVM producers.
  7. Schema evolution: Avro/Protobuf + schema registry from day one — retrofitting from raw JSON is a migration.
  8. Dead-letter topics: route send() callback failures to a local DLQ so retries are bounded; consumer DLQs get retention.ms=2592000000 (30 days).
  9. Partition key strategy: high-cardinality key (order_id, not customer_id) to avoid hot partitions.
  10. Cluster topology: 3 brokers minimum for min.insync.replicas=2 to survive one broker loss; cross-AZ replication in cloud.

Idempotent Producers and Transactional Writes

Two different layers, often conflated: idempotence dedupes retries within one producer session (broker tracks a per-partition sequence number); transactions add atomicity across partitions and sessions. Transactions earn their keep in consume-process-produce — read from topic A, write to topic B, commit the consumer offset atomically. Without them, a crash between produce and offset-commit either loses the downstream record or duplicates it on restart.

Properties props = new Properties();
props.put("bootstrap.servers", "broker1:9092,broker2:9092,broker3:9092");
props.put("transactional.id", "order-processor-instance-7");
props.put("enable.idempotence", "true");
props.put("acks", "all");
props.put("retries", Integer.MAX_VALUE);
props.put("max.in.flight.requests.per.connection", "5");
props.put("delivery.timeout.ms", "120000");
props.put("transaction.timeout.ms", "60000");
 
KafkaProducer<String, OrderEvent> producer = new KafkaProducer<>(props);
producer.initTransactions();  // fences any prior producer with the same transactional.id
 
while (running) {
    ConsumerRecords<String, RawOrder> batch = consumer.poll(Duration.ofMillis(500));
    if (batch.isEmpty()) continue;
 
    producer.beginTransaction();
    try {
        for (ConsumerRecord<String, RawOrder> in : batch) {
            OrderEvent out = transform(in.value());
            producer.send(new ProducerRecord<>("orders.normalized", out.orderId(), out));
        }
        producer.sendOffsetsToTransaction(
            currentOffsets(batch),
            consumer.groupMetadata()
        );
        producer.commitTransaction();
    } catch (ProducerFencedException | OutOfOrderSequenceException fatal) {
        producer.close();
        throw fatal;  // a newer instance has taken over; do not restart this one
    } catch (KafkaException e) {
        producer.abortTransaction();
        // re-poll the same offsets next iteration — consumer has not advanced
    }
}

Two rules: keep transactional.id stable across restarts (hostname + stable ordinal — random UUIDs defeat zombie-fencing), and batch a poll loop's worth of records per transaction, since commitTransaction() blocks ~5–15 ms on the coordinator. Throughput cost: roughly 15–20% on small records, ~5% on large. [Kafka producer config]

Partition Assignment Strategies

Since Kafka 2.4 the default producer partitioner is sticky: it fills one partition's batch until batch.size or linger.ms expires, then rotates — 30–50% better batching on keyless records than legacy round-robin, same topic-level distribution. [Kafka producer config] Keyed records hash with murmur2 mod partition count — deterministic for a fixed count; adding partitions later breaks the mapping, which is why re-partitioning at scale is non-trivial.

For consumers, the assignment strategy is set via partition.assignment.strategy:

# Cooperative sticky — recommended default since Kafka 2.4
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor
 
# Range — assigns contiguous partition ranges per topic; risks skew with multiple topics
# partition.assignment.strategy=org.apache.kafka.clients.consumer.RangeAssignor
 
# Round-robin — distributes evenly but reshuffles all partitions on every rebalance
# partition.assignment.strategy=org.apache.kafka.clients.consumer.RoundRobinAssignor
 
# Sticky (legacy, eager) — minimises movement but stops the world during rebalance
# partition.assignment.strategy=org.apache.kafka.clients.consumer.StickyAssignor
 
session.timeout.ms=45000
heartbeat.interval.ms=3000
max.poll.interval.ms=300000
group.instance.id=consumer-pod-2  # static membership — survives short pod restarts without rebalance

CooperativeStickyAssignor rebalances incrementally — only moving partitions pause, where eager protocols stall the whole group on any pod restart. Custom producer partitioners (tenant isolation, hot-key splitting) implement Partitioner via partitioner.class; keep partition() O(1) and lock-free — it runs on the application thread.

ProducerInterceptor for Auditing

Interceptors attach cross-cutting concerns (audit, schema enforcement, header injection) without touching business code. Both hot-path methods must be fast and non-blocking — onSend() runs on the application thread, onAcknowledgement() on the I/O thread; a slow interceptor stalls every producer in the JVM.

public class AuditingInterceptor implements ProducerInterceptor<String, byte[]> {
 
    private static final Logger AUDIT = LoggerFactory.getLogger("kafka.audit");
    private final AtomicLong sent = new AtomicLong();
    private final AtomicLong failed = new AtomicLong();
 
    @Override
    public ProducerRecord<String, byte[]> onSend(ProducerRecord<String, byte[]> record) {
        record.headers()
            .add("audit.producer", System.getenv("HOSTNAME").getBytes(StandardCharsets.UTF_8))
            .add("audit.timestamp", Long.toString(System.currentTimeMillis()).getBytes())
            .add("audit.trace-id", currentTraceId().getBytes(StandardCharsets.UTF_8));
        return record;
    }
 
    @Override
    public void onAcknowledgement(RecordMetadata metadata, Exception exception) {
        if (exception == null) {
            sent.incrementAndGet();
        } else {
            failed.incrementAndGet();
            AUDIT.warn("send failed: topic={} partition={} cause={}",
                metadata != null ? metadata.topic() : "unknown",
                metadata != null ? metadata.partition() : -1,
                exception.getClass().getSimpleName());
        }
    }
 
    @Override public void close() { /* flush metrics */ }
    @Override public void configure(Map<String, ?> configs) { /* read interceptor config */ }
}

Wire via interceptor.classes (declaration order on send, reverse on ack). Benchmark anything that allocates per record — a 10 µs regression in onSend() is the bottleneck at 100k records/sec.

Kafka 4.0 KRaft-Mode Operational Changes

Kafka 4.0 (early 2026) ships without ZooKeeper — KRaft is mandatory. The producer wire protocol is unchanged; two operational shifts matter to producers: broker config changes now converge in milliseconds via the metadata log (so delivery.timeout.ms can run tighter), and controller failover shrinks from 5–30 s to under a second, so planned controller restarts cause far fewer NOT_CONTROLLER retries.

Frequently Asked Questions

Why enable idempotence?

Before idempotence, retries forced a choice between dropping messages (retries=0) or accepting duplicates. Idempotence adds a per-session sequence number that lets the broker deduplicate retries. The cost is 5 bytes per record — always enable it.

When do I use transactions?

Transactions exist for atomic multi-partition writes and exactly-once consume-process-produce pipelines. They cut throughput by roughly 20%. For single-partition ordering, idempotence plus consistent message keying is enough. [Kafka producer config]

How do I measure whether tuning is working?

Watch record-queue-time-avg (time spent in the producer buffer) and the batch fill ratio (record-size-avg / batch-size-avg). High queue time means linger.ms is too aggressive; a fill ratio below 0.5 means linger.ms is too low and you should raise it.

Keep Reading

Was this article helpful?

Your feedback directly shapes our editorial depth and technical accuracy.

BackendBytes Engineering Team
BackendBytes Engineering Team

Engineering Team

An independent engineering publication covering distributed systems, databases, and production infrastructure. Every factual claim is cited to a primary source or removed.

Read Next