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?

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.

Kafka share consumers

A share consumer reads from a share group where multiple consumers cooperatively consume the same partitions. The broker tracks progress for each individual record rather than by committed offset. This lets you scale the number of consumers beyond the number of partitions and distribute work like a traditional queue.

Share consumers use a separate ShareConsumer class, not the regular Consumer class. The broker drives partition assignment, so there is no rebalance callback and no assign() step. For normal consumers, poll() returns a single record. For share consumers, poll() returns a batch of records that might be empty rather than a single message, so always iterate the result.

Note

Share consumers are a Preview feature. They require confluent-kafka-python 2.15.0 or later. To use share consumers, your Kafka cluster must have share groups enabled, which are available in Confluent Cloud and Confluent Platform 8.2 and later.

A Preview feature is a Confluent component that is introduced to gain early feedback from developers. You can use Preview features for evaluation and non-production testing purposes or to provide feedback to Confluent. The warranty, SLA, and Support Services provisions of your agreement with Confluent don’t apply to Preview features. Confluent might discontinue providing preview releases of the Preview features at any time in Confluent’s sole discretion.

Basic poll loop (implicit acknowledgement)

By default the consumer is in implicit acknowledgement mode: you don’t acknowledge records yourself. Every record returned by a poll is automatically accepted on the next call to poll(), commit_sync(), or commit_async(). Implicit mode is at-least-once, so design your processing to tolerate occasional redelivery.

This example shows a basic share consumer poll loop that uses implicit acknowledgement. For complete examples, see the Python Client repo.

from confluent_kafka import ShareConsumer

consumer = ShareConsumer({
    'bootstrap.servers': 'localhost:9092',
    'group.id': 'my-share-group',
})
consumer.subscribe(['my-topic'])

try:
    while running:
        messages = consumer.poll(timeout=1.0)  # a batch, possibly empty
        for msg in messages:
            if msg.error():
                continue
            process(msg)  # auto-accepted on the next poll
finally:
    consumer.close()

Explicit acknowledgement

Set share.acknowledgement.mode to explicit to acknowledge every record returned by a poll yourself, before the next poll. Acknowledge a record with ACCEPT when processed, RELEASE for transient failures to make it available again for a later delivery attempt, or REJECT for permanent failures, such as a poison record, to stop redelivering it. Flush acknowledgements to the broker with commit_sync(), which blocks and returns a per-partition result, or with commit_async().

from confluent_kafka import ShareConsumer, AcknowledgeType

consumer = ShareConsumer({
    'bootstrap.servers': 'localhost:9092',
    'group.id': 'my-share-group',
    'share.acknowledgement.mode': 'explicit',
})
consumer.subscribe(['my-topic'])

try:
    while running:
        messages = consumer.poll(timeout=1.0)
        for msg in messages:
            if msg.error():
                # librdkafka already acknowledged this record internally
                # based on the error type; no action needed here.
                continue
            try:
                process(msg)
                consumer.acknowledge(msg, AcknowledgeType.ACCEPT)
            except TransientError:
                consumer.acknowledge(msg, AcknowledgeType.RELEASE)
            except Exception:
                consumer.acknowledge(msg, AcknowledgeType.REJECT)
        # Flush acknowledgements. A None value means that partition succeeded.
        results = consumer.commit_sync(timeout=10.0)
        for tp, exc in results.items():
            if exc is not None:
                print(f"commit failed for {tp.topic} [{tp.partition}]: {exc}")
finally:
    consumer.close()

For more information about share consumers, see:

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

3.0.8

FIPS 140-3

3.x

3.1.2

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:

  1. Clone OpenSSL from: OpenSSL Github Repo.

  2. Use git checkout to 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.2 for FIPS 140-3 or git checkout openssl-3.0.8 for FIPS 140-2. The latest OpenSSL version might not be FIPS-compliant.

  3. Run: ./Configure enable-fips.

  4. Run: make install_fips.

    Inside the providers folder, 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() and flush() 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 with statement usage for the producer, consumer, and AdminClient, including automatic resource cleanup when exiting the with block.

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 .