Skip to content

Publishing to Kafka with aiokafka Producers

AIOKafkaProducer batches messages per partition and sends batches in the background; how you call it decides whether that batching happens. Measured with aiokafka 0.14 against a single-node Kafka on the same host, with ~150-byte JSON messages: awaiting send_and_wait for each message, one at a time, managed 2,391 messages per second — every message waited for its own broker round trip. Calling send() for each message and gathering the returned futures managed 152,364 with acks="all", and 158,545 with a 5 ms linger. acks=1 reached 176,342; enabling idempotence cost almost nothing (161,845). When the broker stopped responding, send() still returned immediately, and the failure arrived on the delivery future after the 3.0 s request timeout. This guide sets the producer up for throughput without losing track of what was delivered.

Prerequisites

1. Start one producer per process

A producer holds connections to the brokers, metadata and per-partition buffers. Create one when the process starts and stop it at shutdown — stop() flushes anything still buffered:

from aiokafka import AIOKafkaProducer


@asynccontextmanager
async def lifespan(app):
    producer = AIOKafkaProducer(
        bootstrap_servers="kafka-1:9092,kafka-2:9092",
        acks="all",                    # wait for all in-sync replicas
        enable_idempotence=True,       # no duplicates from producer retries
        linger_ms=5,                   # wait up to 5 ms to fill a batch
        compression_type="zstd",       # or lz4 / gzip; trades CPU for bandwidth
        request_timeout_ms=10_000,
    )
    await producer.start()
    app.state.producer = producer
    yield
    await producer.stop()              # flushes pending batches

A producer per request repeats connection setup and metadata fetches and defeats batching entirely. One producer is safe to share between all tasks in the process; it serializes access internally. If the process forks workers, create the producer in each worker after the fork, for the reasons in why connection pools are per process.

Verify: the broker shows one producer connection per process, stable over time.

2. Send without awaiting each delivery

send() appends the message to a batch and returns a future that resolves when the batch is acknowledged. send_and_wait() does both and waits:

# Slow: one broker round trip per message (2,391 msg/s measured)
for event in events:
    await producer.send_and_wait("orders", value=event, key=event_key)

# Fast: enqueue everything, then wait for all acknowledgements (152,364 msg/s)
futures = [await producer.send("orders", value=e, key=k) for e, k in batch]
results = await asyncio.gather(*futures)       # RecordMetadata: partition, offset

The await on send() itself is short — it only waits if the producer's buffer is full, which is the producer's backpressure. The delivery futures are what matter: gathering them confirms every message was acknowledged, and each result carries the partition and offset. Never drop the futures and move on; an unobserved failure is a lost message with no error anywhere.

Verify: the gathered results include a partition and offset for every message sent.

Messages per second by producer usage 6 horizontal bars comparing send_and_wait per message with the others. Messages per second by producer usage send_and_wait per message 2,391/s send + gather, acks=all 152,364/s + linger_ms=5 158,545/s + gzip 140,221/s idempotent, linger 5 161,845/s acks=1, linger 5 176,342/s aiokafka 0.14, single-node Kafka on localhost, 3 partitions, keyed messages; 20,000 messages (2,000 for the first row). Not waiting per message is worth 60x; the settings after that are fine-tuning.

3. Choose acks and idempotence for durability

acks decides when the broker acknowledges a batch, and therefore what a resolved delivery future promises:

AIOKafkaProducer(acks=0)              # no acknowledgement: fastest, can lose anything
AIOKafkaProducer(acks=1)              # leader wrote it: lost if the leader dies before replicating
AIOKafkaProducer(acks="all")          # all in-sync replicas wrote it: survives a broker loss
AIOKafkaProducer(acks="all", enable_idempotence=True)   # and producer retries cannot duplicate

Measured on one node, acks=1 was about 11% faster than acks="all"; on a real cluster the gap is larger because all waits for replication across the network. For events that matter — orders, payments, state changes — use acks="all" with enable_idempotence=True, and set the topic's min.insync.replicas to 2 so "all" means at least two copies. Idempotence makes the broker discard duplicate batches when the producer retries after a lost acknowledgement; it measured at no meaningful cost (161,845 against 158,545 msg/s).

Verify: the topic has min.insync.replicas ≥ 2, and the producer config shows acks=all with idempotence for business events.

4. Use keys for ordering and spread

The message key chooses the partition. Messages with the same key go to the same partition and are read in order; messages without keys are spread across partitions:

await producer.send("orders", value=payload, key=str(order.customer_id).encode())

Key by the entity whose events must stay in order — a customer, an account, an order id — so a consumer sees that entity's events in sequence. Avoid low-cardinality keys (a country code, a status) that put most traffic on one partition, and remember that changing the partition count remaps keys to partitions, which breaks per-key ordering across the change. Consumers then process each partition in order and can parallelize across keys, as in processing queue items in order per key.

Verify: per-partition message rates are within a reasonable range of each other, and all events for one key land in one partition.

The life of a produced message A flow of 5 stages. The life of a produced message send(key, value) append to batch linger / batch full batch sent leader + replicas per acks setting future resolves partition, offset or fails after request timeout The future is the only proof of delivery; always observe it.

5. Handle delivery failures

When the broker is unreachable or slow, send() still returns at once — the failure comes later, on the future. Tested by pausing the broker: send() returned in under a millisecond, and the future failed with RequestTimedOutError after 3.01 s (the configured request_timeout_ms). After the broker resumed, the next send succeeded normally:

from aiokafka.errors import KafkaError


async def publish(producer, topic: str, key: bytes, value: bytes) -> None:
    future = await producer.send(topic, value=value, key=key)
    try:
        await future
    except KafkaError:
        log.exception("delivery failed for key %s", key)
        await outbox.save(topic, key, value)       # retry later from durable storage
        raise

What to do with a failed message depends on what produced it. Inside a request handler, return an error so the caller can retry. For events that must not be lost, write them to an outbox table in the same database transaction as the state change, and publish from the outbox — the transactional outbox pattern — so a broker outage delays events instead of losing them. Monitor delivery failures as a metric; a producer that fails silently is the most common way events go missing.

Verify: with the broker paused in a test, failures surface as errors or outbox rows, and every message is delivered once the broker returns.

How should this producer be configured? A decision on What is being published with 3 outcomes. How should this producer be configured? What is being published? business events acks=all + idempotence outbox on failure high-volume telemetry acks=1, linger, compression loss tolerable needs per-entity order key by entity one partition per key Throughput comes from not waiting per message; durability from acks.

Verification

A producer is set up well when:

  • One producer per process starts at startup and stops (flushing) at shutdown.
  • Messages are sent without per-message waits, and every delivery future is observed.
  • Business events use acks="all" with idempotence and adequate min.insync.replicas.
  • Delivery failures are surfaced or saved, never ignored.

Diagnostic Hook: export the producer's send rate, delivery failure count and the time between send() and future resolution. Resolution time rising toward request_timeout_ms means the brokers are struggling; failures without any corresponding rise mean misconfiguration — an unknown topic, an authorization error — rather than load.

Pitfalls & edge cases

  • send_and_wait in a loop. Measured at 2,391 msg/s against 152,364 with gathered futures.
  • Discarding delivery futures. Failures disappear without a trace.
  • A producer per request. No batching and repeated connection setup.
  • Low-cardinality keys. One hot partition limits throughput and consumers.

Frequently Asked Questions

How do I send Kafka messages quickly with aiokafka?

Call await producer.send() for each message, which only enqueues it, then await the returned futures together. In testing that gave 152,364 messages per second, against 2,391 when awaiting send_and_wait for each message.

Should I use acks=all with aiokafka?

For events that must not be lost, yes, together with enable_idempotence=True and a topic min.insync.replicas of at least 2. acks=1 was only about 11% faster on a single-node test.

What happens when Kafka is down and I call producer.send()?

send() returns immediately; the returned future fails after request_timeout_ms. In testing it raised RequestTimedOutError after 3.0 s. Always await or check the future.

Do I need to close an aiokafka producer?

Yes. await producer.stop() at shutdown flushes buffered batches; without it, messages still in the buffer are lost.