Python Client for Apache Kafka
Python Client overview
Confluent offers the confluent-kafka-python client for Apache Kafka®. The client
provides a high-level producer, consumer, and AdminClient that are
compatible with Kafka brokers version 0.8 and later, Confluent Cloud, and Confluent Platform. For
release updates, see Kafka Python Client Changelog.
This page shows how to install the client, configure and use the producer and consumer, use share consumers (Preview), configure FIPS-compliant communication, and use the AsyncIO producer and consumer.
Ready to get started?
Sign up for Confluent Cloud, the fully managed cloud-native service for Apache Kafka® and get started for free using the Cloud quick start.
Download Confluent Platform, the self managed, enterprise-grade distribution of Apache Kafka and get started using the Confluent Platform quick start.
The confluent-kafka-python package is a binding on top of the C client, librdkafka. For an overview of the librdkafka client library, see Introduction to librdkafka client library.
For information about the configuration of the Confluent Kafka Python client, see Kafka client configuration.
Installation
The client is available on PyPI and can be installed using pip:
pip install confluent-kafka
You can install it globally, or within a virtualenv. If you want to install a FIPS-compliant client, see FIPS compliance.
Note
The confluent-kafka-python package comes bundled with a pre-built version of librdkafka which does not include GSSAPI/Kerberos support. For information about how to install a version that supports GSSAPI, see the installation instructions.
Tutorials and sample code
For a step-by-step tutorial that uses the Python client, including code samples for the producer and consumer, see this guide.
For examples using basic producers, consumers, AsyncIO, and how to produce and consume Avro data with Schema Registry, see confluent-kafka-python GitHub repository.
Kafka producer
Initialize the producer
The producer is configured using a dictionary in the examples below. If you are running Kafka locally, you can initialize the producer as shown below.
Tip
You can use the Python with context manager to automatically close the producer and consumer objects.
For a working example, see context_manager_example.py
from confluent_kafka import Producer
import socket
conf = {'bootstrap.servers': 'host1:9092,host2:9092',
'client.id': socket.gethostname()}
producer = Producer(conf)
If you are connecting to a Kafka cluster in Confluent Cloud, you must provide credentials for access. The example below shows using a cluster API key and secret.
from confluent_kafka import Producer
import socket
conf = {'bootstrap.servers': 'pkc-abcd85.us-west-2.aws.confluent.cloud:9092',
'security.protocol': 'SASL_SSL',
'sasl.mechanism': 'PLAIN',
'sasl.username': '<CLUSTER_API_KEY>',
'sasl.password': '<CLUSTER_API_SECRET>',
'client.id': socket.gethostname()}
producer = Producer(conf)
For information on the available configuration properties, see the API Documentation.
Asynchronous writes
To initiate sending a message to Kafka, call the produce method, passing in the
message value (which might be None) and optionally a key, partition, and callback.
The produce call completes immediately and does not return a value. A
BufferError is raised if the message could not be enqueued due to
librdkafka’s local produce queue being full.
producer.produce(topic, key="key", value="value")
To receive notification of delivery success or failure, you can pass a callback
parameter. This can be any callable, for example, a lambda, function, bound method, or
callable object. Although the produce() method enqueues message immediately
for batching, compression and transmission to broker, no delivery notification events
are propagated until poll() is invoked.
def acked(err, msg):
if err is not None:
print("Failed to deliver message: %s: %s" % (str(msg), str(err)))
else:
print("Message produced: %s" % (str(msg)))
producer.produce(topic, key="key", value="value", callback=acked)
# Wait up to one second for events. Callbacks will be invoked during
# this method call if the message is acknowledged.
producer.poll(1)
Synchronous writes
The Python client provides a flush() method which can be used to make
writes synchronous. This is typically a bad idea since it effectively
limits throughput to the broker round trip time, but might be justified in
some cases.
producer.produce(topic, key="key", value="value")
producer.flush()
Typically, flush() should be called prior to shutting down the producer
to ensure all outstanding/queued/in-flight messages are delivered.
Kafka consumer
Initialize the consumer
The consumer is configured using a dictionary in the examples below. If you are running Kafka locally, you can initialize the consumer as shown below.
Tip
You can use the Python with context manager to automatically close the producer and consumer objects.
For a working example, see context_manager_example.py
from confluent_kafka import Consumer
conf = {'bootstrap.servers': 'host1:9092,host2:9092',
'group.id': 'foo',
'auto.offset.reset': 'earliest'}
consumer = Consumer(conf)
If you are connecting to a Kafka cluster in Confluent Cloud, you must provide credentials for access. The example below shows using a cluster API key and secret.
from confluent_kafka import Consumer
conf = {'bootstrap.servers': 'pkc-abcd85.us-west-2.aws.confluent.cloud:9092',
'security.protocol': 'SASL_SSL',
'sasl.mechanism': 'PLAIN',
'sasl.username': '<CLUSTER_API_KEY>',
'sasl.password': '<CLUSTER_API_SECRET>',
'group.id': 'foo',
'auto.offset.reset': 'earliest'}
consumer = Consumer(conf)
The group.id property is mandatory and specifies which consumer group the consumer
is a member of. The auto.offset.reset property specifies what offset the consumer
should start reading from in the event there are no committed offsets for a partition,
or the committed offset is invalid (perhaps due to log truncation).
For information on the available configuration properties, see the API Documentation.
Consumer code examples
Basic poll loop
A typical Kafka consumer application is centered around a consume loop, which repeatedly calls
the poll method to retrieve records one-by-one that have been efficiently pre-fetched by
the consumer in behind the scenes. Before entering the consume loop, you’ll typically use the
subscribe method to specify which topics should be fetched from:
from confluent_kafka import KafkaException
running = True
def basic_consume_loop(consumer, topics):
try:
consumer.subscribe(topics)
while running:
msg = consumer.poll(timeout=1.0)
if msg is None: continue
if msg.error():
raise KafkaException(msg.error())
else:
msg_process(msg)
finally:
# Close down consumer to commit final offsets.
consumer.close()
def shutdown():
global running
running = False
The poll timeout is hard-coded to one second. If no records
are received before this timeout expires, then Consumer.poll() returns
None.
Note that you should always call Consumer.close() after you are finished
using the consumer. Doing so ensures that active sockets are
closed and internal state is cleaned up. It also triggers a group
rebalance immediately which ensures that any partitions owned by the
consumer are re-assigned to another member in the group. If not closed
properly, the broker triggers the rebalance only after the session
timeout has expired.
Synchronous commits
To commit offsets manually, first disable automatic commits by setting
enable.auto.commit to false:
from confluent_kafka import Consumer
conf = {'bootstrap.servers': 'host1:9092,host2:9092',
'group.id': 'foo',
'enable.auto.commit': 'false',
'auto.offset.reset': 'earliest'}
consumer = Consumer(conf)
The simplest and most reliable way to manually commit offsets is by
setting the asynchronous parameter of the Consumer.commit()
method call to False. This method can also accept the mutually exclusive keyword
parameters offsets to explicitly list the offsets for each assigned
topic partition and message which commits offsets relative to a
Message object returned by poll().
def consume_loop(consumer, topics):
try:
consumer.subscribe(topics)
msg_count = 0
while running:
msg = consumer.poll(timeout=1.0)
if msg is None: continue
if msg.error():
raise KafkaException(msg.error())
else:
msg_process(msg)
msg_count += 1
if msg_count % MIN_COMMIT_COUNT == 0:
consumer.commit(asynchronous=False)
finally:
# Close down consumer to commit final offsets.
consumer.close()
In this example, a synchronous commit is triggered every MIN_COMMIT_COUNT
messages. The asynchronous flag controls whether this call is
asynchronous. You could also trigger the commit on expiration of a
timeout to ensure there the committed position is updated regularly.
Asynchronous commits
Asynchronous commits send the commit request and return immediately without waiting for the broker to acknowledge it, which avoids blocking the poll loop.
def consume_loop(consumer, topics):
try:
consumer.subscribe(topics)
msg_count = 0
while running:
msg = consumer.poll(timeout=1.0)
if msg is None: continue
if msg.error():
raise KafkaException(msg.error())
else:
msg_process(msg)
msg_count += 1
if msg_count % MIN_COMMIT_COUNT == 0:
consumer.commit(asynchronous=True)
finally:
# Commit final offsets synchronously; close() doesn't commit
# when enable.auto.commit is false.
consumer.commit(asynchronous=False)
consumer.close()
In this example, the consumer sends the request and returns
immediately by using asynchronous commits. The example sets the
asynchronous parameter of commit() to True explicitly, but
asynchronous commits are the default if you omit the parameter.
The API invokes a commit callback when the commit either succeeds or
fails. The callback can be any callable. Pass it to the consumer
constructor as the on_commit configuration property.
from confluent_kafka import Consumer
def commit_completed(err, partitions):
if err:
print(str(err))
else:
print("Committed partition offsets: " + str(partitions))
conf = {'bootstrap.servers': "host1:9092,host2:9092",
'group.id': "foo",
'auto.offset.reset': 'earliest',
'on_commit': commit_completed}
consumer = Consumer(conf)
Delivery guarantees
How and when the application stores and commits offsets determines the delivery semantics. This section covers the default behavior, which provides no delivery guarantee, and how to get “at least once” or “at most once” delivery.
Default guarantee
The following are the default settings of the Python client, and the settings result in the default delivery guarantee of the Python client being “None”:
'enable.auto.commit': 'true'
'enable.auto.offset.store': 'true'
Because auto commits run in a background thread, these settings might result in the offset for the latest message being committed before the application has finished processing the message. If the application were to crash or exit before finishing processing, and the offset had been auto-committed, the next incarnation of the consumer application would start at the next message, effectively missing the message that was processed when the application crashed. You can lose data or get duplicates.
“At least once” guarantee
To achieve “at least once” guarantee, configure the following settings:
'enable.auto.commit': 'true'
'enable.auto.offset.store': 'false'
To avoid data loss with the default configuration
(enable.auto.commit: true, enable.auto.offset.store: true),
the application can disable the automatic offset store and manually store
offsets with store_offsets() after processing. Messages might be
redelivered after a failure, so design processing to tolerate duplicates.
This gives the application fine-grained control over when a message is
committed. The latest stored offset is automatically committed every
auto.commit.interval.ms.
For this guarantee option, you should store the offset only after processing the message successfully.
def consume_loop(consumer, topics):
try:
consumer.subscribe(topics)
while running:
msg = consumer.poll(timeout=1.0)
if msg is None: continue
if msg.error():
raise KafkaException(msg.error())
else:
msg_process(msg)
consumer.store_offsets(message=msg)
finally:
# Close down consumer to commit final offsets.
consumer.close()
Note
Only offsets greater than the current offset are committed. For example,
if the latest committed offset was 10 and the application calls
store_offsets() with offset 9, that offset is not committed.
“At most once” guarantee
Committing synchronously before processing the message, instead of after, gives “at most once” delivery: the offset is committed even if processing later fails, so a crash can skip a message rather than reprocess it. You must handle commit failures carefully: if a commit fails and you process the message anyway, a later redelivery can process it twice, which breaks the “at most once” guarantee.
def consume_loop(consumer, topics):
try:
consumer.subscribe(topics)
while running:
msg = consumer.poll(timeout=1.0)
if msg is None: continue
if msg.error():
raise KafkaException(msg.error())
else:
consumer.commit(message=msg, asynchronous=False)
msg_process(msg)
finally:
# Close down consumer to commit final offsets.
consumer.close()
For simplicity, this example calls Consumer.commit() before
processing each message. Committing on every message would
produce a lot of overhead in practice. A better approach would be to
collect a batch of messages, execute the synchronous commit, and then
process the messages only if the commit succeeded.
FIPS compliance
This client supports both FIPS 140-2 and FIPS 140-3 compliance. Use the following version mapping to ensure compliance:
Compliance Standard |
OpenSSL Version |
FIPS Provider Version |
|---|---|---|
FIPS 140-2 |
3.x |
|
FIPS 140-3 |
3.x |
For new deployments, use FIPS 140-3. On September 21, 2026, all remaining FIPS 140-2 certificates moved to the CMVP historical list, so federal procurement no longer accepts FIPS 140-2 modules for new acquisitions.
Confluent tested communication for FIPS compliance between clients and the following endpoints:
Kafka brokers
Schema Registry
Kafka broker and Schema Registry
- Kafka broker
To communicate with the Kafka broker, the client uses the librdkafka library, which uses OpenSSL. The steps below configure OpenSSL to operate in FIPS mode for client communication with the broker.
- Schema Registry
To communicate with Schema Registry, the client uses the standard libraries in Python, which uses the operating system’s native SSL/TLS library. For FIPS-compliant communication between Schema Registry and the client, do not use the steps that follow. Instead, make the SSL/TLS library FIPS compliant. If the native SSL/TLS library is OpenSSL (the default for Python), then use the steps in the OpenSSL readme to make OpenSSL FIPS compliant.
FIPS-compliant Kafka broker communication
For FIPS-compliant communication with the broker, you can install the client in two ways:
Use prebuilt wheels
Build librdkafka and the client both from source
Use prebuilt wheels
If you install this client through prebuilt wheels using pip install confluent_kafka, OpenSSL 3.0
is already statically linked with the librdkafka shared library. To enable this client to
communicate with the Kafka cluster using the OpenSSL FIPS provider and FIPS-approved algorithms,
you must enable the FIPS provider. You can find steps to enable the FIPS provider in section
Use FIPS provider.
Note
You should enable the FIPS provider (using the same steps) if you install this client from the source using
pip install confluent_kafka --no-binary :all: with prebuilt librdkafka in which OpenSSL is statically linked.
When the client uses a statically linked OpenSSL, such as with prebuilt
wheels, set both the OPENSSL_MODULES and OPENSSL_CONF environment
variables in the environment where the client runs.
Build librdkafka and client both from source
When you build the librdkafka from source, librdkafka dynamically links to the OpenSSL present in the system
if static linking is not used explicitly while building. If the system installed OpenSSL is already working in FIPS mode,
then you can directly jump to the section
client configuration to enable FIPS provider and enable the fips provider.
If you don’t have OpenSSL working in FIPS mode, use the steps mentioned in the section
Use FIPS provider to make OpenSSL in your system FIPS compliant, and then enable the fips
provider. After you have OpenSSL working in FIPS mode and the fips provider enabled, librdkafka and the Python Client
use FIPS-approved algorithms to communicate with the Kafka cluster.
Use FIPS provider
To use the FIPS provider, you must have the FIPS module available on your system. Plug the module into OpenSSL, and then configure OpenSSL to use the module.
Plug the FIPS provider into OpenSSL as described in Reference FIPS provider in OpenSSL.
After you plug the FIPS provider module into OpenSSL, configure OpenSSL to use the module. You have two options:
Edit the default configuration file to include the FIPS-related configuration.
Create a new configuration file and point to it using the environment variable,
OPENSSL_CONF. For example:OPENSSL_CONF="/path/to/fips/enabled/openssl/config/openssl.cnf.For an example of OpenSSL configuration file, see: Enable FIPS provider with OpenSSL.
Build FIPS provider module
Building the FIPS provider module means cloning OpenSSL, checking out a FIPS-compliant version, and running the build steps below. For the official steps to generate the FIPS provider module, see the following version-specific documentation:
To build the FIPS provider module:
Clone OpenSSL from: OpenSSL Github Repo.
Use
git checkoutto check out the OpenSSL tag listed in the FIPS Provider Version column of the version mapping table in FIPS compliance for your compliance standard. For example,git checkout openssl-3.1.2for FIPS 140-3 orgit checkout openssl-3.0.8for FIPS 140-2. The latest OpenSSL version might not be FIPS-compliant.Run:
./Configure enable-fips.Run:
make install_fips.Inside the
providersfolder, two files are generated. Use these files with OpenSSL:FIPS module (fips.dylib in Mac, fips.so in Linux, and fips.dll in Windows)
FIPS configuration file:
fipsmodule.cnf
Reference FIPS provider in OpenSSL
You can plug the FIPS provider into OpenSSL in two ways: place the module in
the default OpenSSL module folder, or point to the folder that contains it with
the OPENSSL_MODULES environment variable (for example,
OPENSSL_MODULES="/path/to/fips/module/lib/folder/"). The default module
folder works only when librdkafka uses a dynamically linked OpenSSL. When
OpenSSL is statically linked, such as with prebuilt wheels, use
OPENSSL_MODULES.
When installing from source, plug the FIPS provider module you built
(fips.so, fips.dylib, or fips.dll) into OpenSSL by placing it in
the default OpenSSL module folder. Look for something like:
...lib/ossl-modules/.
For the default locations of OpenSSL on various operating systems, see the SSL section of the Introduction to librdkafka - the Apache Kafka C/C++ client library.
You can also point to this module with the environment variable OPENSSL_MODULES.
For example: OPENSSL_MODULES="/path/to/fips/module/lib/folder/.
Enable FIPS provider with OpenSSL
To enable FIPS in OpenSSL, you must include fipsmodule.cnf in the file, openssl.cnf. See the
following openssl.cnf example:
config_diagnostics = 1
openssl_conf = openssl_init
.include /usr/local/ssl/fipsmodule.cnf
[openssl_init]
providers = provider_sect
alg_section = algorithm_sect
[provider_sect]
fips = fips_sect
[algorithm_sect]
default_properties = fips=yes
.
.
.
The fipsmodule.cnf file includes fips_sect which OpenSSL requires to enable FIPS.
Some of the algorithms might have different implementation in FIPS or other providers. If you load two different
providers like default and fips, any implementation could be used. To make sure you fetch only FIPS-compliant version
of the algorithm, use fips=yes default property in the configuration file.
Client configuration to enable FIPS provider
To make a client FIPS compliant, set 'ssl.providers': 'fips,base' in its
configuration. This enables both the fips provider and the base
provider, which OpenSSL requires for non-crypto algorithms that the FIPS
provider doesn’t include. The base provider ships with OpenSSL by default.
AsyncIO support
The Python client provides AsyncIO-compatible producer and consumer clients for
integration with asynchronous Python applications. Use the AsyncIO clients when your
application uses Python’s asyncio framework and you require non-blocking Kafka
operations that integrate seamlessly with other asynchronous operations.
Note
The AsyncIO classes are available in the confluent_kafka.aio package in
confluent-kafka-python 2.13.0 and later. Earlier versions provide them as
experimental classes in confluent_kafka.experimental.aio.
When to use AsyncIO clients
Choose the AsyncIO or synchronous client based on your application architecture and the criteria below.
Use AsyncIO clients when:
Your application runs under an event loop (FastAPI, Starlette, aiohttp, Sanic, asyncio workers).
You must avoid blocking the event loop during Kafka operations.
You must integrate Kafka operations with other asynchronous I/O operations.
You’re building asynchronous web services or microservices.
Use synchronous clients when:
Building scripts, batch jobs, or CLI tools.
Running high-throughput data pipelines where you control threads and processes.
You are able to call
poll()andflush()directly without blocking issues.You require per-message headers in produce operations (not supported in asynchronous batched path).
In an asynchronous server environment, prefer using the AsyncIO client for
better compatibility. If you require headers, invoke the synchronous produce() using
run_in_executor() for that path.
AsyncIO producer
Initialize the AsyncIO producer
Import AIOProducer from confluent_kafka.aio and configure it using a
dictionary, similar to the synchronous Producer.
For Confluent Cloud connections, provide credentials in the configuration dictionary:
import asyncio
from confluent_kafka.aio import AIOProducer
async def produce_to_ccloud():
conf = {
'bootstrap.servers': 'pkc-abcd85.us-west-2.aws.confluent.cloud:9092',
'security.protocol': 'SASL_SSL',
'sasl.mechanism': 'PLAIN',
'sasl.username': '<CLUSTER_API_KEY>',
'sasl.password': '<CLUSTER_API_SECRET>',
'client.id': 'my-async-producer'
}
producer = AIOProducer(conf)
try:
# Your async produce logic here
pass
finally:
await producer.close()
asyncio.run(produce_to_ccloud())
For Confluent Platform connections, configure the AIOProducer as follows.
This example assumes a local, unsecured cluster.
import asyncio
from confluent_kafka.aio import AIOProducer
async def produce_example():
producer = AIOProducer({
'bootstrap.servers': 'localhost:9092',
'client.id': 'my-async-producer'
})
try:
# Your async produce logic here
pass
finally:
await producer.close()
asyncio.run(produce_example())
Produce messages asynchronously
Call the produce method with await to enqueue messages for delivery. The
method returns a Future that resolves to the delivered message.
import asyncio
async def send_message(producer, topic):
# Produce returns a Future
delivery_future = await producer.produce(
topic=topic,
key='key1',
value='value1'
)
# Await the future to get delivery confirmation
msg = await delivery_future
print(f'Delivered to {msg.topic()} [{msg.partition()}] @ {msg.offset()}')
asyncio.run(send_message(producer, topic))
Produce messages in batches
Produce multiple messages concurrently and await their delivery using
asyncio.gather():
import asyncio
async def batch_produce(producer, topic):
# Create multiple produce futures
futures = [
await producer.produce(topic=topic, key=f'key{i}', value=f'value{i}')
for i in range(100)
]
# Flush to ensure messages are in flight
await producer.flush()
# Wait for all deliveries
messages = await asyncio.gather(*futures)
print(f'Delivered {len(messages)} messages')
asyncio.run(batch_produce())
Produce messages transactionally
The AIOProducer supports asynchronous transactional operations:
import asyncio
from confluent_kafka.aio import AIOProducer
async def transactional_produce(topic):
producer = AIOProducer({
'bootstrap.servers': 'host1:9092,host2:9092',
'transactional.id': 'my-transactional-producer'
})
await producer.init_transactions()
try:
await producer.begin_transaction()
# Produce messages within transaction
futures = [
await producer.produce(topic=topic, value=f'msg{i}')
for i in range(10)
]
await producer.flush()
await asyncio.gather(*futures)
# Commit the transaction
await producer.commit_transaction()
except Exception as e:
# Abort on error
await producer.abort_transaction()
raise
finally:
await producer.close()
asyncio.run(transactional_produce('my-topic'))
Close the producer
Call await producer.close() to release the producer’s resources.
close() flushes any buffered messages and waits for delivery
confirmation before returning, so pending messages aren’t dropped.
AsyncIO consumer
Initialize the AsyncIO consumer
Import AIOConsumer from confluent_kafka.aio and configure it
using a dictionary.
For Confluent Cloud connections, provide credentials in the configuration dictionary:
import asyncio
from confluent_kafka.aio import AIOConsumer
async def consume_from_ccloud():
conf = {
'bootstrap.servers': 'pkc-abcd85.us-west-2.aws.confluent.cloud:9092',
'security.protocol': 'SASL_SSL',
'sasl.mechanism': 'PLAIN',
'sasl.username': '<CLUSTER_API_KEY>',
'sasl.password': '<CLUSTER_API_SECRET>',
'group.id': 'my-consumer-group',
'auto.offset.reset': 'earliest'
}
consumer = AIOConsumer(conf)
try:
# Your async consume logic here
pass
finally:
await consumer.close()
asyncio.run(consume_from_ccloud())
For Confluent Platform connections, configure the AIOConsumer as follows.
This example assumes a local, unsecured cluster.
import asyncio
from confluent_kafka.aio import AIOConsumer
async def consume_example():
consumer = AIOConsumer({
'bootstrap.servers': 'localhost:9092',
'group.id': 'my-consumer-group',
'auto.offset.reset': 'earliest'
})
try:
# Your async consume logic here
pass
finally:
await consumer.close()
asyncio.run(consume_example())
Consume messages in a loop
Use await with the poll() method to consume messages without blocking
the event loop:
import asyncio
from confluent_kafka.aio import AIOConsumer
# Use the same configuration keys shown in the initialization section
kafka_configuration_object = {
'bootstrap.servers': 'localhost:9092',
'group.id': 'my-consumer-group',
'auto.offset.reset': 'earliest'
}
topics = ['my-topic']
async def main():
consumer = AIOConsumer(kafka_configuration_object)
await consumer.subscribe(topics)
try:
while True:
# Poll without blocking the event loop
msg = await consumer.poll(timeout=1.0)
if msg is None:
continue
if msg.error():
print(f'Consumer error: {msg.error()}')
continue
# Process the message
print(f'Received: {msg.value().decode("utf-8")}')
finally:
await consumer.unsubscribe()
await consumer.close()
asyncio.run(main())
Manage offsets manually
Disable auto.commit and manually commit offsets for precise control:
import asyncio
from confluent_kafka.aio import AIOConsumer
kafka_configuration_object = {
'bootstrap.servers': 'localhost:9092',
'group.id': 'my-consumer-group',
'enable.auto.commit': 'false',
'enable.auto.offset.store': 'false',
'auto.offset.reset': 'earliest'
}
topics = ['my-topic']
async def main():
consumer = AIOConsumer(kafka_configuration_object)
await consumer.subscribe(topics)
msg_count = 0
try:
while True:
msg = await consumer.poll(timeout=1.0)
if msg is None:
continue
if msg.error():
continue
# Process the message
print(f'Received: {msg.value()}')
# Store offset after processing
await consumer.store_offsets(message=msg)
msg_count += 1
# Commit every 100 messages
if msg_count % 100 == 0:
await consumer.commit()
finally:
await consumer.unsubscribe()
await consumer.close()
asyncio.run(main())
Perform rebalance with asynchronous callbacks
Define asynchronous callback functions for partition assignment and revocation:
import asyncio
from confluent_kafka.aio import AIOConsumer
kafka_configuration_object = {
'bootstrap.servers': 'localhost:9092',
'group.id': 'my-consumer-group',
'enable.auto.offset.store': 'false',
'auto.offset.reset': 'earliest'
}
topics = ['my-topic']
async def main():
consumer = AIOConsumer(kafka_configuration_object)
async def on_assign(consumer, partitions):
print(f'Partitions assigned: {partitions}')
# If this callback doesn't call assign() (eager assignors) or
# incremental_assign() (cooperative assignors) itself, the client
# assigns these partitions automatically after it returns.
async def on_revoke(consumer, partitions):
print(f'Partitions revoked: {partitions}')
await consumer.commit()
async def on_lost(consumer, partitions):
print(f'Partitions lost: {partitions}')
await consumer.subscribe(
topics,
on_assign=on_assign,
on_revoke=on_revoke,
on_lost=on_lost
)
try:
while True:
msg = await consumer.poll(timeout=1.0)
if msg is not None and not msg.error():
print(f'Consumed: {msg.value()}')
await consumer.store_offsets(message=msg)
finally:
await consumer.unsubscribe()
await consumer.close()
asyncio.run(main())
Integrate with Schema Registry
Integrate with Schema Registry by using the asynchronous AsyncSchemaRegistryClient
and the matching asynchronous serializer or deserializer class (for example,
AsyncAvroSerializer), shown in the examples below. These asynchronous Schema Registry
classes are experimental: their APIs, including the initialization pattern,
might change in future versions. The asynchronous serializer and deserializer
classes are imported from private modules (for example, _async) to signal
their unstable status.
Produce Avro messages asynchronously
Serialize a value with an asynchronous Avro serializer, then await
producer.produce() to send it, as shown below.
import asyncio
from confluent_kafka.schema_registry import AsyncSchemaRegistryClient
from confluent_kafka.schema_registry._async.avro import AsyncAvroSerializer
from confluent_kafka.serialization import SerializationContext, MessageField
from confluent_kafka.aio import AIOProducer
async def produce_avro_async():
# Configure async Schema Registry client
sr_client = AsyncSchemaRegistryClient({
'url': 'http://localhost:8081'
})
# Define Avro schema...
schema_str = '''
{
"type": "record",
"name": "User",
"fields": [{"name": "name", "type": "string"}]
}
'''
# Await the serializer constructor to complete async initialization
avro_serializer = await AsyncAvroSerializer(
sr_client,
schema_str=schema_str
)
producer = AIOProducer({'bootstrap.servers': 'localhost:9092'})
try:
# Serialize and produce
value = {'name': 'alice'}
serialized = await avro_serializer(
value,
SerializationContext('my-topic', MessageField.VALUE)
)
delivery_future = await producer.produce(
'my-topic',
value=serialized
)
msg = await delivery_future
print(f'Delivered to {msg.topic()}')
finally:
await producer.flush()
await producer.close()
asyncio.run(produce_avro_async())
Consume Avro messages asynchronously
Deserialize a consumed message’s value with an asynchronous Avro deserializer after polling, as shown below.
import asyncio
from confluent_kafka.schema_registry import AsyncSchemaRegistryClient
from confluent_kafka.schema_registry._async.avro import AsyncAvroDeserializer
from confluent_kafka.aio import AIOConsumer
from confluent_kafka.serialization import SerializationContext, MessageField
async def consume_avro_async():
# Configure async Schema Registry client
sr_client = AsyncSchemaRegistryClient({
'url': 'http://localhost:8081'
})
# Await the deserializer constructor to complete async initialization
avro_deserializer = await AsyncAvroDeserializer(sr_client)
consumer = AIOConsumer({
'bootstrap.servers': 'localhost:9092',
'group.id': 'my-avro-consumer',
'auto.offset.reset': 'earliest'
})
await consumer.subscribe(['my-topic'])
try:
msg = await consumer.poll(timeout=5.0)
if msg is not None and not msg.error():
value = await avro_deserializer(
msg.value(),
SerializationContext('my-topic', MessageField.VALUE)
)
print(f'Received Avro value: {value}')
finally:
await consumer.unsubscribe()
await consumer.close()
asyncio.run(consume_avro_async())
Examples
Producer and consumer examples
The basic producer and consumer examples work with any Kafka deployment, but are optimized for Confluent Cloud and Confluent Platform which provide additional features, security, and enterprise support.
producer.py Read lines from stdin and send them to a Kafka topic.
consumer.py Read messages from a Kafka topic.
context_manager_example.py Demonstrates context manager
withstatement usage for the producer, consumer, and AdminClient, including automatic resource cleanup when exiting thewithblock.
AsyncIO examples
The AsyncIO examples examples demonstrate AsyncIO patterns including concurrent producer and consumer operations, signal handling, and transaction management.
asyncio_example.py Comprehensive AsyncIO example that demonstrates both AIOProducer and AIOConsumer with transactional operations, batched asynchronous produce, proper event loop integration, signal handling, and asynchronous callback patterns.
asyncio_avro_producer.py Minimal AsyncIO Avro producer using AsyncSchemaRegistryClient and AsyncAvroSerializer. Supports Confluent Cloud using
--sr-api-key/--sr-api-secret.
API documentation
To view the Python client API documentation, click here .