Kafka — Producers & Consumers in Python¶
All examples use confluent-kafka 2.15.1 (librdkafka 2.15.1) against a Kafka 4.3.1 broker from 02 Local Setup & CLI.
Python Clients¶
| Package | Version checked | What it is | When to use |
|---|---|---|---|
confluent-kafka |
2.15.1 | Wrapper over librdkafka (C), maintained by Confluent; wheels include librdkafka | Default choice: fastest, feature-complete (transactions, KIP-848, admin API), Schema Registry serializers, asyncio API in confluent_kafka.aio |
kafka-python |
3.0.11 | Pure-Python client, releases resumed after a long pause | No native dependency allowed; simple scripts |
aiokafka |
0.14.0 | Pure-Python asyncio client (aio-libs) | Existing asyncio code bases already built on it |
uv add confluent-kafka # producer, consumer, admin
uv add "confluent-kafka[avro,schemaregistry]" # + Avro serializers and Schema Registry client
Config keys are librdkafka keys
confluent-kafka takes a plain dict with the dotted names from librdkafka's CONFIGURATION.md (bootstrap.servers, enable.idempotence, …). Most match the Java client; some defaults do not — see the tables below.
Admin: Creating Topics from Code¶
from confluent_kafka.admin import AdminClient, NewTopic
admin = AdminClient({"bootstrap.servers": "localhost:9092"})
futures = admin.create_topics([
NewTopic("orders", num_partitions=3, replication_factor=1, config={"retention.ms": "604800000"}),
NewTopic("orders.dlq", num_partitions=1, replication_factor=1),
])
for topic, future in futures.items():
try:
future.result() # raises KafkaException, e.g. TOPIC_ALREADY_EXISTS
print(f"created {topic}")
except Exception as exc:
print(f"{topic}: {exc}")
print(sorted(t for t in admin.list_topics(timeout=5).topics if not t.startswith("__")))
Admin calls are asynchronous and return a dict of futures — always call .result(), or errors are silently lost.
Producer Configuration¶
| Setting | librdkafka default | Recommended for business events | Notes |
|---|---|---|---|
acks |
all |
all |
1 loses data if the leader dies before followers copy it |
enable.idempotence |
false |
true |
No duplicates or reordering on retries; forces acks=all and at most 5 in-flight requests. Java clients default to true |
retries |
2147483647 |
default | Retries are bounded by time, not count |
delivery.timeout.ms |
300000 (5 min) |
30–120 s | Total time to deliver one message including retries; after it the delivery report fails |
linger.ms |
5 |
5–50 | Wait to fill a batch: more throughput, a bit more latency |
batch.size |
1000000 bytes |
default | Max batch size in bytes |
compression.type |
none |
zstd or lz4 |
Compresses whole batches; brokers store them compressed |
partitioner |
consistent_random (CRC32) |
murmur2_random if Java producers write the same topic |
Same key, same partition across languages |
client.id |
rdkafka |
service name + host | Shows up in broker logs, quotas, metrics |
Producing with Delivery Reports¶
import json
import socket
from confluent_kafka import KafkaError, Message, Producer
producer = Producer({
"bootstrap.servers": "localhost:9092",
"client.id": f"orders-api-{socket.gethostname()}",
"acks": "all",
"enable.idempotence": True,
"linger.ms": 5,
"compression.type": "zstd",
"delivery.timeout.ms": 30000,
})
def on_delivery(err: KafkaError | None, msg: Message) -> None:
if err is not None:
print(f"FAILED {msg.key()}: {err}") # real code: log, metric, alert
else:
print(f"ok {msg.key().decode()} -> {msg.topic()}[{msg.partition()}]@{msg.offset()}")
for order_id, status in [("ord-1", "created"), ("ord-2", "created"), ("ord-1", "paid")]:
producer.produce(
"orders",
key=order_id, # str or bytes; str is UTF-8 encoded
value=json.dumps({"order_id": order_id, "status": status}).encode(),
headers={"event-type": "order-status-changed", "schema-version": "1"},
on_delivery=on_delivery,
)
producer.poll(0) # serve callbacks of earlier messages
remaining = producer.flush(10) # wait for all outstanding deliveries
assert remaining == 0, f"{remaining} messages not delivered"
How the producer works:
produce()only puts the message into a local queue and returns immediately. It does not mean the broker has the message.- A background thread batches and sends; the result arrives as a delivery report.
- Delivery callbacks run only inside
poll()orflush()— callpoll(0)in the produce loop,flush()before exit. flush(timeout)returns the number of messages still undelivered; a non-zero value means data can be lost on exit.
| Situation | What happens | Handle it |
|---|---|---|
Local queue full (queue.buffering.max.messages, default 100000) |
produce() raises BufferError: Local: Queue full |
producer.poll(1) to drain, then retry the same produce() |
| Topic does not exist and auto-creation is off | Message waits in the queue; after delivery.timeout.ms the report fails with _MSG_TIMED_OUT, and a shorter flush() returns > 0 |
Create topics in setup; check the flush() result and delivery reports |
| Broker down | librdkafka retries and reconnects until delivery.timeout.ms |
Alert on failed delivery reports |
Message larger than message.max.bytes (1000000) |
produce() raises KafkaException with MSG_SIZE_TOO_LARGE |
Store the payload elsewhere, send a reference |
Headers¶
- Send as a dict or a list of
(key, value)tuples; values can bestrorbytes. - A consumer gets
msg.headers()as a list of(str, bytes)tuples, orNonewhen there are no headers. - Typical uses:
event-type,schema-version,correlation-id,traceparent(added by OpenTelemetry instrumentation), error details on DLQ messages.
Consumer Configuration¶
| Setting | librdkafka default | Notes |
|---|---|---|
group.id |
— (required for subscribe) |
One group per logical consumer service |
auto.offset.reset |
latest |
Where to start when the group has no committed offset (or it is out of range). earliest reads the backlog; latest only new records |
enable.auto.commit |
true |
Commits stored offsets in the background every auto.commit.interval.ms (5000) |
enable.auto.offset.store |
true |
Marks a message as done as soon as poll() returns it |
isolation.level |
read_committed |
Skips records of aborted transactions |
max.poll.interval.ms |
300000 |
Max time between poll() calls before the member is kicked out of the group |
session.timeout.ms |
45000 |
Heartbeat timeout (classic protocol) |
group.protocol |
classic |
consumer = KIP-848 protocol |
group.instance.id |
— | Static membership: a restarted pod with the same ID gets its partitions back without a rebalance |
auto.offset.reset only applies to a group without offsets
Once a group has committed, it resumes from there and auto.offset.reset is ignored. A test that reuses group.id="tests" and "sometimes sees nothing" is usually reading from an old committed offset — use a new group per test.
Commit Strategies¶
| Strategy | Config | Guarantee | Throughput |
|---|---|---|---|
| Auto commit, auto store (defaults) | — | A crash after poll() but before processing finishes can skip messages |
High |
| Auto commit, manual store | enable.auto.offset.store=false, store_offsets(msg) after processing |
At-least-once | High |
| Sync commit per message | enable.auto.commit=false, commit(message=msg, asynchronous=False) |
At-least-once, smallest duplicate window | Low: one round trip per message |
| Commit every N messages / T seconds | enable.auto.commit=false, periodic commit(asynchronous=False) |
At-least-once, up to N duplicates after a crash | High |
| Transactions | transactional.id on the producer |
Exactly-once for Kafka-to-Kafka | Medium |
At-Least-Once Consumer Loop¶
import json
import signal
from confluent_kafka import Consumer, KafkaException
running = True
def stop(signum, frame):
global running
running = False
signal.signal(signal.SIGTERM, stop) # Kubernetes / docker stop
consumer = Consumer({
"bootstrap.servers": "localhost:9092",
"group.id": "billing",
"auto.offset.reset": "earliest",
"enable.auto.commit": True,
"enable.auto.offset.store": False, # we decide what is done
})
consumer.subscribe(["orders"])
def handle(event: dict) -> None:
print("processing", event)
try:
while running:
msg = consumer.poll(1.0)
if msg is None:
continue
if msg.error():
raise KafkaException(msg.error())
handle(json.loads(msg.value()))
consumer.store_offsets(message=msg) # committed by the next auto commit
finally:
consumer.close() # commits stored offsets and leaves the group
close()matters: it commits what was stored and leaves the group, so its partitions are reassigned at once instead of aftersession.timeout.ms.- A handler that can take longer than
max.poll.interval.msgets the consumer kicked out mid-work. Split the work, raise the limit, orpause()the partitions and keep callingpoll(). - The consumer is meant to be used from one thread. The producer is thread-safe and can be shared.
Rebalancing¶
from confluent_kafka import Consumer
def on_assign(consumer, partitions):
print("assigned", [(p.topic, p.partition) for p in partitions])
def on_revoke(consumer, partitions):
print("revoked", [(p.topic, p.partition) for p in partitions])
# manual commits: commit processed offsets here, before another member takes over
consumer = Consumer({
"bootstrap.servers": "localhost:9092",
"group.id": "billing",
"group.protocol": "consumer", # KIP-848: broker-side, incremental assignment
"auto.offset.reset": "earliest",
"enable.auto.commit": False,
})
consumer.subscribe(["orders"], on_assign=on_assign, on_revoke=on_revoke)
- With
group.protocol=consumer,session.timeout.ms,heartbeat.interval.msandpartition.assignment.strategyare rejected by the client (_INVALID_ARG ... It is defined broker side); usegroup.remote.assignorto pick the server-side assignor. - On the classic protocol set
partition.assignment.strategy=cooperative-stickyto avoid stop-the-world rebalances. kafka-groups.sh --listshows each group's type (ClassicorConsumer).
Exactly-Once: Consume-Transform-Produce¶
import json
from confluent_kafka import Consumer, KafkaException, Producer
consumer = Consumer({
"bootstrap.servers": "localhost:9092",
"group.id": "invoicer",
"auto.offset.reset": "earliest",
"enable.auto.commit": False,
"isolation.level": "read_committed",
})
producer = Producer({
"bootstrap.servers": "localhost:9092",
"transactional.id": "invoicer-1", # stable per instance; enables idempotence
})
producer.init_transactions()
consumer.subscribe(["orders"])
while True:
msgs = consumer.consume(num_messages=100, timeout=1.0)
if not msgs:
continue
producer.begin_transaction()
try:
for msg in msgs:
if msg.error():
raise KafkaException(msg.error())
order = json.loads(msg.value())
producer.produce("invoices", key=msg.key(), value=json.dumps({"invoice_for": order["order_id"]}).encode())
producer.send_offsets_to_transaction( # input offsets commit together with the output
consumer.position(consumer.assignment()),
consumer.consumer_group_metadata(),
)
producer.commit_transaction()
except Exception:
producer.abort_transaction()
raise # real code: seek back to the committed offsets and retry
- Output records and input offsets are committed atomically: after a crash, either both are visible or neither is.
- Downstream consumers must use
read_committed(the librdkafka default) or they will see aborted records. - Every transaction writes a commit marker into each partition it touched, which takes an offset. In a check run, 3 messages in one transaction moved the log-end offset to 4 — do not assert
log-end offset == message counton transactional topics. - It does not cover external side effects: a DB write inside the loop can still happen twice.
asyncio: confluent_kafka.aio¶
import asyncio
from confluent_kafka.aio import AIOConsumer, AIOProducer
async def main() -> None:
producer = AIOProducer({"bootstrap.servers": "localhost:9092"})
try:
delivery = await producer.produce("orders", key="ord-9", value=b'{"order_id": "ord-9"}')
msg = await delivery # resolves on the delivery report
print("delivered", msg.partition(), msg.offset())
finally:
await producer.flush()
await producer.close()
consumer = AIOConsumer({
"bootstrap.servers": "localhost:9092",
"group.id": "aio-demo",
"auto.offset.reset": "earliest",
})
try:
await consumer.subscribe(["orders"])
for _ in range(10):
msg = await consumer.poll(1.0)
if msg is not None and not msg.error():
print("consumed", msg.key(), msg.offset())
finally:
await consumer.close()
asyncio.run(main())
AIOProducer / AIOConsumer run the blocking librdkafka calls in a thread pool and expose the same methods as coroutines — handy inside FastAPI or other asyncio services.
Common Mistakes¶
| Mistake | Effect | Fix |
|---|---|---|
No flush() before a script or test exits |
Last messages never sent | assert producer.flush(10) == 0 |
| Ignoring delivery reports | Failed writes go unnoticed | on_delivery callback that logs and counts errors |
| Relying on defaults for business events | No idempotence: retries can duplicate or reorder | enable.idempotence=True |
| Auto commit + auto store with slow handlers | Offsets committed before processing ends; a crash skips messages | Manual store or manual commit after processing |
| Consumer not closed | Partitions blocked until the session times out; stored offsets not committed | try/finally: consumer.close() |
| Long blocking work between polls | Member kicked out, rebalance loop, duplicates | Keep work under max.poll.interval.ms, or pause() / resume() |
| New producer per message | Connection setup per call, very slow | One long-lived producer per process |
See also¶
- Apache Kafka — Overview
- Kafka — Core Concepts
- Kafka — Schemas & Observability
- Kafka — Testing Kafka-Based Systems
- Queues vs Streams: Message Delivery, Ordering & Reliability
- FastAPI — Modern Async Web Framework