Part 3 · 1 chapters · ~8 min
Producers
The producer pipeline (partitioner, record accumulator, sender), batching with batch.size and linger.ms, compression, in-flight requests, idempotence with producer ids and sequence numbers, retries and delivery timeouts, keys and ordering, and handling send failures.
5
Batching, compression and idempotence
code
// KafkaJS-style producer configuration (concepts are the same in every client)
const producer = kafka.producer({ idempotent: true, maxInFlightRequests: 5 });
await producer.send({
topic: 'transfers', acks: -1, compression: CompressionTypes.ZSTD, // acks=all
messages: [{ key: accountId, value: JSON.stringify(event), headers: { 'event-id': event.id, 'traceparent': tp } }],
});
# Java/librdkafka equivalents: enable.idempotence=true, acks=all, linger.ms=10, batch.size=65536, compression.type=zstdINSIDE A PRODUCER
records batched per partition, compressed, sent in flight, retried idempotently
swipe the figure sideways, or tap expand for full screen
1/5
batching
send() is asynchronous: records are added to a batch per partition in the accumulator. Batches are sent when full (batch.size) or after linger.ms. Small linger values (5-20 ms) raise throughput a lot.
batches per partitionlinger.ms trades a little latency for throughput