NewYour coding agent can read the release notes before it upgrades.Set up the MCP server →
PyPI · #471 most downloaded on PyPI
Confluent's Python client for Apache Kafka
Last release 6 days ago
28 Sep 2026
Ships fairly regularly
a new release about every 5 weeks
Nearly every release is documented
notes for 46 of 51 stable releases
2 versions withdrawn
withdrawn after publishing
10 years old
67 releases · first in 2016
On Windows, the tag isn't currently given when building librdkafka, t…
On Windows, the tag isn't currently given when building librdkafka, t…
v2.16.0 is a feature release with the following features, fixes and enhancements:
confluent_kafka now declares itself GIL-safe, enabling real multi-core parallelism on free-threaded CPython builds. See the Multithreading Guide for thread-safety details and free-threaded caveats. (#2347)close() now aborts any open transaction (#2347)DlqAction). When a rule fails, the record is teed to a configured DLQ
topic and the original serialize/deserialize call still raises. With the
default (global) RuleRegistry the DLQ is best-effort; set
dlq.auto.flush=true or give the serde its own RuleRegistry (closable on
shutdown) for durability.Producer, Consumer, AdminClient, and
Message classes (#2347)Consumer instance across threads
instead of leaving it as undefined behavior (#2347)Producer.purge() ignoring in_queue, in_flight and blocking set to False on big-endian platforms such as s390x (#2345)One column per quarter.
2.16.0 RC1 Syntax and style fixes
2.16.0 RC1
Syntax and style fixes
v2.16.0 is a feature release with the following features, fixes and enhancements:
DlqAction). When a rule fails, the record is teed to a configured DLQ
topic and the original serialize/deserialize call still raises. With the
default (global) RuleRegistry the DLQ is best-effort; set
dlq.auto.flush=true or give the serde its own RuleRegistry (closable on
shutdown) for durability.DeserializingConsumer and DeserializingShareConsumer now deserialize the message key before the value, matching SerializingProducer and the Java clien
DeserializingConsumer and DeserializingShareConsumer now deserialize theSerializingProducer and the JavaAdminClient.delete_records() followedlist_topics()) on Python 3.14.asyncio.get_running_loop() instead of asyncio.get_event_loop() to avoid creating a new event loop and raise an error in case a loop isn't available (AlexCai26, #2339).confluent-kafka-python 2.15.1 is based on librdkafka 2.15.1, see the
librdkafka release notes
for a complete list of changes, enhancements, fixes and upgrade considerations.
2.15.1rc2: make validation pass on setup_all_versions.py
2.15.1rc2: make validation pass on setup_all_versions.py (#2354)
KIP-932 Queues for Kafka – Now in Preview
confluent-kafka-python 2.15.0 adds a Preview implementation of the KIP-932 share consumer (Queues for Kafka). Members of a share group cooperatively consume from the same partitions with per-record acquire/acknowledge semantics and redelivery, providing queue-like consumption on top of Kafka.
ShareConsumer and DeserializingShareConsumer clients: subscribe, batch poll (poll() returns a Messages batch), and close (with context-manager support).share.acknowledgement.mode), with AcknowledgeType ACCEPT / RELEASE / REJECT and Message.delivery_count().commit_sync / commit_async) and an acknowledgement-commit callback.set_sasl_credentials) and new IllegalStateException / ConcurrentModificationException exceptions.See the Share Consumer guide and the Share consumers section of librdkafka's INTRODUCTION.md.
Note: The share consumer is currently in Preview and should not be used in production environments. Its public interfaces may change before General Availability, and known limitations apply (see the guide). The share consumer is single-threaded and not thread-safe. It requires a broker with share groups enabled (generally available in Apache Kafka 4.2.0).
confluent-kafka[oauthbearer-aws] provides AWS IAM-basedGetWebIdentityToken. Activate by setting sasl.oauthbearer.method=oidc,sasl.oauthbearer.metadata.authentication.type=aws_iam, andsasl.oauthbearer.config="region=...,audience=...". See theexamples/oauth_oidc_ccloud_aws_iam.pyconfluent-kafka-python v2.15.0 is based on librdkafka v2.15.0, see the
librdkafka release notes
for a complete list of changes, enhancements, fixes and upgrade considerations.
v2.14.2 is a maintenance release with the following fixes and enhancements:
v2.14.2 is a maintenance release with the following fixes and enhancements:
confluent-kafka-python v2.14.2 is based on librdkafka v2.14.2, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
Nothing published for this version
v2.14.0 is a feature release with the following features, fixes and enhancements:
v2.14.0 is a feature release with the following features, fixes and enhancements:
confluent-kafka-python v2.14.0 is based on librdkafka v2.14.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
Nothing published for this version
v2.13.2 is a maintenance release with the following fixes and enhancements:
v2.13.2 is a maintenance release with the following fixes and enhancements:
confluent-kafka-python v2.13.2 is based on librdkafka v2.13.2, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
pip install confluent-kafka==2.13.2
Nothing published for this version
v2.13.0 is a feature release with the following features, fixes and enhancements:
v2.13.0 is a feature release with the following features, fixes and enhancements:
close() method to producer (#2039)__len__ function to AIOProducer (#2140)__enter__ to return the same object type that called it (#2157)Consumer.poll(), Consumer.consume(), Producer.poll(), and Producer.flush() blocking indefinitely and not responding to Ctrl+C (KeyboardInterrupt) signals. The implementation now uses a "wakeable poll" pattern that breaks long blocking calls into smaller chunks (200ms) and periodically re-acquires the Python GIL to check for pending signals. This allows Ctrl+C to properly interrupt blocking operations. Fixes Issues #209 and #807 (#2126).list_consumer_group_offsets() (#2118)commit() and store_offsets() (#2145)confluent-kafka-python v2.13.0 is based on librdkafka v2.13.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
pip install confluent-kafka==2.13.0
This is a beta release intended for --pre installers to get early access to upcoming changes. Exact changes included are:
This is a beta release intended for --pre installers to get early access to upcoming changes. Exact changes included are:
close() method to producerlist_consumer_group_offsets()v2.12.2 is a maintenance release with the following change:
v2.12.2 is a maintenance release with the following change:
Accept-Version to Confluent-Accept-Unknown-Properties in request headerv2.12.2 is a hotfix for a critical problem found with Schema Registry clients in the 2.12.1 release:
v2.12.1 is a maintenance release with the following fixes:
v2.12.1 is a maintenance release with the following fixes:
libversion() now returns the string/integer tuple for varients on version -- use version() for string only responsesr.lookup_schema()Nothing published for this version
There are few contract change associated with the new protocol and might cause breaking changes. group.protocol configuration property dictates whethe…
confluent-kafka-python v2.12.0
v2.12.0 is a feature release with the following enhancements:
Starting with confluent-kafka-python 2.12.0, the next generation consumer group rebalance protocol defined in KIP-848 is production-ready. Please refer to the following migration guide for moving from classic to consumer protocol.
Note: The new consumer group protocol defined in KIP-848 is not enabled by default. There are few contract change associated with the new protocol and might cause breaking changes. group.protocol configuration property dictates whether to use the new consumer protocol or older classic protocol. It defaults to classic if not provided.
Introduces beta class AIOProducer for asynchronous message production in asyncio applications.
AIOProducer for
asynchronous message production in asyncio applications. This API offloads
blocking librdkafka calls to a thread pool and schedules common callbacks
(error_cb, throttle_cb, stats_cb, oauth_cb, logger) onto the event
loop for safe usage inside async frameworks.await AIOProducer(...).produce(topic, value=...)
buffers messages and flushes when the buffer threshold or timeout is reached.await producer.flush(), await producer.purge(), and
transactional operations (init_transactions, begin_transaction,
commit_transaction, abort_transaction).Producer.produce(...) or
offload a sync produce call to a thread executor within your async app.Producer remains recommended.confluent-kafka-python v2.12.0 is based on librdkafka v2.12.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
Nothing published for this version
Nothing published for this version
confluent-kafka-python v2.11.1 is based on librdkafka v2.11.1, see the librdkafka release notes for a complete list of changes, enhancements, fixes an
confluent-kafka-python v2.11.1 is based on librdkafka v2.11.1, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
confluent-kafka-python v2.11.0 is based on librdkafka v2.11.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes an
confluent-kafka-python v2.11.0 is based on librdkafka v2.11.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
v2.11.0 is a feature release with the following enhancements:
confluent-kafka-python v2.11.0 is based on librdkafka v2.11.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
c56a3e6
This commit was created on GitHub.com and signed with GitHub’s verified signature .
GPG key ID: B5690EEEBB952194
Verified Learn about vigilant mode .
librdkafka v2.11.0 is a feature release:
KIP-1102 Enable clients to rebootstrap based on timeout or error code ( #4981 ).
KIP-1139 Add support for OAuth jwt-bearer grant type ( #4978 ).
Fix for poll ratio calculation in case the queues are forwarded ( #5017 ).
Fix data race when buffer queues are being reset instead of being initialized ( #4718 ).
Features BROKER_BALANCED_CONSUMER and SASL_GSSAPI don't depend on JoinGroup v0 anymore, missing in AK 4.0 and CP 8.0 ( #5131 ).
Improve HTTPS CA certificates configuration by probing several paths when OpenSSL is statically linked and providing a way to customize their location or value ( #5133 ).
Nothing published for this version
Handled None value for optional ctx parameter in ProtobufDeserializer
None value for optional ctx parameter in ProtobufDeserializer (#1939)None value for optional ctx parameter in AvroDeserializer (#1973)confluent-kafka-python v2.10.1 is based on librdkafka v2.10.1, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
v2.10.1 is a maintenance release with the following fixes:
None value for optional ctx parameter in ProtobufDeserializer (#1939)None value for optional ctx parameter in AvroDeserializer (#1973)confluent-kafka-python v2.10.1 is based on librdkafka v2.10.1, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
Nothing published for this version
[KIP-848] Group Config is now supported in AlterConfigs, IncrementalAlterConfigs and DescribeConfigs.
describe_consumer_groups() now supports KIP-848 introduced consumer groups. Two new fields for consumer group type and target assignment have also been added. Type defines whether this group is a classic or consumer group. Target assignment is only valid for the consumer protocol and its defaults to NULL. (#1873).confluent-kafka-python v2.10.0 is based on librdkafka v2.10.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
Add Client Credentials OAuth support for Schema Registry (#1919) Add custom OAuth support for Schema Registry (#1925) confluent-kafka-python v2.9.0 is
Add Client Credentials OAuth support for Schema Registry (#1919) Add custom OAuth support for Schema Registry (#1925) confluent-kafka-python v2.9.0 is based on librdkafka v2.8.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
v2.9.0 is a feature release with the following fixes and enhancements:
confluent-kafka-python v2.9.0 is based on librdkafka v2.8.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
Note: Versioning is skipped due to breaking change in v2.8.1. Do not run software with v2.8.1 installed.
confluent-kafka-python v2.8.2 is based on librdkafka v2.8.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
Note: Versioning is skipped due to breaking change in v2.8.1. Do not run software with v2.8.1 installed.
Nothing published for this version
Ensure algorithm query param is passed for CSFLE
confluent-kafka-python v2.8.0 is based on librdkafka v2.8.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
Note: As part of this release, we are deprecating v2.6.2 release and yanking it from PyPI. Please refrain from using v2.6.2. Use v2.7.0 instead.
Note: As part of this release, we are deprecating v2.6.2 release and yanking it from PyPI. Please refrain from using v2.6.2. Use v2.7.0 instead.
Note: This release modifies the dependencies of the Schema Registry client. If you are using the Schema Registry client, please ensure that you install the extra dependencies using the following syntax:
pip install confluent-kafka[schemaregistry]
or
pip install confluent-kafka[avro,schemaregistry]
Please see the README.md for more information related to installing protobuf, jsonschema or rules dependencies.
confluent-kafka-python v2.7.0 is based on librdkafka v2.6.1, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
v2.7.0 is a feature release with the features, fixes and enhancements present in v2.6.2 including the following fix:
confluent-kafka-python v2.7.0 is based on librdkafka v2.6.1, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
Upon the release of 2.7.0, the 2.6.2 version will be marked deprecated. We apologize for the inconvenience and appreciate the feedback that we have go…
[!WARNING] Due to an error in which we included dependency changes to a recent patch release, Confluent recommends users to refrain from upgrading to 2.6.2 of Confluent Kafka. Confluent will release a new minor version, 2.7.0, where the dependency changes will be appropriately included. Users who have already upgraded to 2.6.2 and made the required dependency changes are free to remain on that version and are recommended to upgrade to 2.7.0 when that version is available. Upon the release of 2.7.0, the 2.6.2 version will be marked deprecated. We apologize for the inconvenience and appreciate the feedback that we have gotten from the community.
Note: This version is yanked from PyPI. Use 2.7.0 instead.
Note: This release modifies the dependencies of the Schema Registry client. If you are using the Schema Registry client, please ensure that you install the extra dependencies using the following syntax:
pip install confluent-kafka[schemaregistry]
or
pip install confluent-kafka[avro,schemaregistry]
Please see the README.md for more information related to installing protobuf, jsonschema or rules dependencies.
confluent-kafka-python is based on librdkafka v2.6.1, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
Migrated build system from setup.py to pyproject.toml in accordance with PEP 517 and PEP 518, improving project configuration, build system requiremen
setup.py to pyproject.toml in accordance with PEP 517 and PEP 518, improving project configuration, build system requirements management, and compatibility with modern Python packaging tools like pip and build. (#1592)confluent-kafka-python is based on librdkafka v2.6.1, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
Added Python 3.13 wheels (#1828).
confluent-kafka-python is based on librdkafka v2.6.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
Release asset checksums:
v2.11.0.zip SHA256 9e76a408f0ed346f21be5e2df58b672d07ff9c561a5027f16780d1b26ef24683
v2.11.0.tar.gz SHA256 592a823dc7c09ad4ded1bc8f700da6d4e0c88ffaf267815c6f25e7450b9395ca
Fix an assert being triggered during push telemetry call when no metrics matched on the client side.
confluent-kafka-python is based on librdkafka v2.5.3, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
> [!WARNING] This version has introduced a regression in which an assert is triggered during PushTelemetry call. This happens when no metric is matche
[!WARNING] This version has introduced a regression in which an assert is triggered during PushTelemetry call. This happens when no metric is matched on the client side among those requested by broker subscription.
You won't face any problem if:
- Broker doesn't support KIP-714.
- KIP-714 feature is disabled on the broker side.
- KIP-714 feature is disabled on the client side. This is enabled by default. Set configuration
enable.metrics.pushtofalse.- If KIP-714 is enabled on the broker side and there is no subscription configured there.
- If KIP-714 is enabled on the broker side with subscriptions that match the KIP-714 metrics defined on the client.
Having said this, we strongly recommend using
v2.5.3and above to not face this regression at all.
AdminClient. (#1758)strcpy to enhance security of the client. (#1745)OAUTHBEARER/OIDC extensions copy. (#1745)operation_timeout and request_timeout in various Admin apis. (#1710)TopicCollection and TopicPartitionInfo classes when importing through other module like mypy. (#1764)commit or store_offsets consumer method is called incorrectly with errored Message object. (#1754)logger not working when provided as an argument to AdminClient instead of a configuration property. (#1758)PyDict_SetItem. (#1710)confluent-kafka-python is based on librdkafka v2.5.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
KIP-848: Added KIP-848 based new consumer group rebalance protocol. The feature is an Early Access: not production ready yet. Please refer detailed do
confluent-kafka-python is based on librdkafka v2.4.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
KIP-117: Add support for AdminAPI describe_cluster() and describe_topics(). (@jainruchir, #1635)
describe_cluster() and describe_topics(). (@jainruchir, #1635)list_offsets (#1576).Rack to the Node type, so AdminAPI calls can expose racks for brokers
(currently, all Describe Responses) (#1635, @jainruchir).confluent-kafka-python is based on librdkafka v2.3.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
KIP-339 IncrementalAlterConfigs API (#1517).
confluent-kafka-python is based on librdkafka v2.2.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
Added a new ConsumerGroupState UNKNOWN. The typo state UNKOWN is deprecated and will be removed in the next major version.
confluent-kafka-python is based on librdkafka v2.1.1, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
Added set_sasl_credentials. This new method (on the Producer, Consumer, and AdminClient) allows modifying the stored SASL PLAIN/SCRAM credentials that
set_sasl_credentials. This new method (on the Producer, Consumer, and AdminClient) allows modifying the stored SASL PLAIN/SCRAM credentials that will be used for subsequent (new) connections to a broker (#1511).confluent-kafka-python is based on librdkafka v2.1.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
Deprecated AvroProducer and AvroConsumer. Use AvroSerializer and AvroDeserializer instead.
list_consumer_groups Admin operation. Supports listing by state.describe_consumer_groups Admin operation. Supports multiple groups.delete_consumer_groups Admin operation. Supports multiple groups.list_consumer_group_offsets Admin operation. Currently, only supports 1 group with multiple partitions. Supports require_stable option.alter_consumer_group_offsets Admin operation. Currently, only supports 1 group with multiple offsets.normalize.schemas configuration property to Schema Registry client (@rayokota, #1406)TopicPartition type and commit() (#1410).consumer.memberid() for getting member id assigned to
the consumer in a consumer group (#1154).nb_bool method for the Producer, so that the default (which uses len)
will not be used. This avoids situations where producers with no enqueued items would
evaluate to False (@vladz-sternum, #1445).AvroProducer and AvroConsumer. Use AvroSerializer and AvroDeserializer instead.list_groups. Use list_consumer_groups and describe_consumer_groups instead.OpenSSL 3.0.x upgrade in librdkafka requires a major version bump, as some legacy ciphers need to be explicitly configured to continue working, but it is highly recommended NOT to use them. The rest of the API remains backward compatible.
confluent-kafka-python is based on librdkafka v2.0.2, see the librdkafka v2.0.0 release notes and later ones for a complete list of changes, enhancements, fixes and upgrade considerations.
Note: There were no v2.0.0 and v2.0.1 releases.
Support for setting principal and SASL extensions in oauth_cb and handle failures (@Manicben, #1402)
confluent-kafka-python is based on librdkafka v1.9.2, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
The warnings for use.deprecated.format (introduced in v1.8.2) had its logic reversed, which result in warning logs to be emitted when the property was…
confluent-kafka-python is based on librdkafka v1.9.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
IMPORTANT: Added mandatory use.deprecated.format to ProtobufSerializer and ProtobufDeserializer. See Upgrade considerations below for more information…
v1.8.2 is a maintenance release with the following fixes and enhancements:
use.deprecated.format to ProtobufSerializer and ProtobufDeserializer.
See Upgrade considerations below for more information.use.latest.version and skip.known.types (Protobuf) to the Serializer classes. (Robert Yokota, #1133).list_topics() and list_groups() added to AdminClient.confluent-kafka-python is based on librdkafka v1.8.2, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
Note: There were no v1.8.0 and v1.8.1 releases.
Prior to this version the confluent-kafka-python client had a bug where nested protobuf schemas indexes were incorrectly serialized, causing incompatibility with other Schema-Registry protobuf consumers and producers.
This has now been fixed, but since the old defect serialization and the new correct serialization are mutually incompatible the user of confluent-kafka-python will need to make an explicit choice which serialization format to use during a transitory phase while old producers and consumers are upgraded.
The ProtobufSerializer and ProtobufDeserializer constructors now both take a (for the time being) configuration dictionary that requires
the use.deprecated.format configuration property to be explicitly set.
Producers should be upgraded first and as long as there are old (<=v1.7.0) Python consumers reading from topics being produced to, the new (>=v1.8.2) Python producer must be configured with use.deprecated.format set to True.
When all existing messages in the topic have been consumed by older consumers the consumers should be upgraded and both new producers and the new consumers must set use.deprecated.format to False.
The requirement to explicitly set use.deprecated.format will be removed in a future version and the setting will then default to False (new format).
v1.7.0 is a maintenance release with the following fixes and enhancements:
v1.7.0 is a maintenance release with the following fixes and enhancements:
confluent-kafka-python is based on librdkafka v1.7.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
Fix deprecated schema.Parse call (@casperlehmann, #1006).
v1.6.1 is a feature release:
return_record_name=True to AvroDeserializer (@slominskir, #1028)schema.Parse call (@casperlehmann, #1006).**kwargs to legacy AvroProducer and AvroConsumer constructors to
support all Consumer and Producer base class constructor arguments, such
as logger (@venthur, #699).producer.flush() could return a non-zero value without hitting the specified timeout.confluent-kafka-python is based on librdkafka v1.6.1, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
v1.6.0 is a feature release with the following features, fixes and enhancements:
v1.6.0 is a feature release with the following features, fixes and enhancements:
Message.latency() to retrieve the per-message produce latency.consumer.close() must now be explicitly called if the application
wants to leave the consumer group properly and commit final offsets.PY_SSIZE_T_CLEAN warningproducer.purge() to purge messages in-queue/flight (@peteryin21, #548)AdminClient.list_groups() API (@messense, #948)confluent-kafka-python is based on librdkafka v1.6.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
v1.5.0 is a maintenance release with the following fixes and enhancements:
v1.5.0 is a maintenance release with the following fixes and enhancements:
confluent-kafka-python is based on librdkafka v1.5.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
v1.4.2 is a maintenance release with the following fixes and enhancements:
v1.4.2 is a maintenance release with the following fixes and enhancements:
confluent-kafka-python is based on librdkafka v1.4.2, see the librdkafka v1.4.2 release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
Nothing published for this version
KIP-98: Transactional Producer API
v1.4.0 is a feature release:
confluent-kafka-python is based on librdkafka v1.4.0, see the librdkafka v1.4.0 release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
Release v1.4.0 for confluent-kafka-python adds complete Exactly-Once-Semantics (EOS) functionality, supporting the idempotent producer (since v1.0.0), a transaction-aware consumer (since v1.2.0) and full producer transaction support (v1.4.0).
This enables developers to create Exactly-Once applications with Apache Kafka.
See the Transactions in Apache Kafka page for an introduction and check the transactions example.
Release v1.4.0 introduces a new, experimental, API which adds serialization capabilities to Kafka Producer and Consumer. This feature provides the ability to configure Producer/Consumer key and value serializers/deserializers independently. Previously all serialization must be handled prior to calling Producer.produce and after Consumer.poll.
This release ships with 3 built-in, Java compatible, standard serializer and deserializer classes:
| Name | Type | Format |
|---|---|---|
| Double | float | IEEE 764 binary64 |
| Integer | int | int32 |
| String | Unicode | bytes* |
* The StringSerializer codec is configurable and supports any one of Python's standard encodings. If left unspecified 'UTF-8' will be used.
Additional serialization implementations are possible through the extension of the Serializer and Deserializer base classes.
See avro_producer.py and avro_consumer.py for example usage.
Release v1.4.0 for confluent-kafka-python adds support for two new Schema Registry serialization formats with its Generic Serialization API; JSON and Protobuf. A new set of Avro Serialization classes have also been added to conform to the new API.
| Format | Serializer Example | Deserializer Example |
|---|---|---|
| Avro | avro_producer.py | avro_consumer.py |
| JSON | json_producer.py | json_consumer.py |
| Protobuf | protobuf_producer.py | protobuf_consumer.py |
Two security issues have been identified in the SASL SCRAM protocol handler:
sasl.username and sasl.password contained characters that needed escaping, a buffer overflow and heap corruption would occur. This was protected, but too late, by an assertion.Both of these issues are fixed in this release.
General:
Schema Registry/Avro:
/ from Schema Registry base URL (@IvanProdaiko94 , #749)Also see the librdkafka v1.4.0 release notes for fixes to the underlying client implementation.
Upgrade builtin lz4 to 1.9.2 (CVE-2019-17543, #2598).
confluent-kafka-python is based on librdkafka v1.3.0, see the librdkafka v1.3.0 release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
This is a feature release adding support for KIP-392 Fetch from follower, allowing a consumer to fetch messages from the closest replica to increase throughput and reduce cost.
confluent-kafka-python is based on librdkafka v1.2.0, see the librdkafka v1.2.0 release notes for a complete list of changes, enhancements, fixes and
confluent-kafka-python is based on librdkafka v1.2.0, see the librdkafka v1.2.0 release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
isolation.level=read_committed) implemented by @mhowlett.linger.ms) on the producer.This release adds consumer-side support for transactions.
In previous releases, the consumer always delivered all messages to the application, even those in aborted or not yet committed transactions. In this release, the consumer will by default skip messages in aborted transactions.
This is controlled through the new isolation.level configuration property which
defaults to read_committed (only read committed messages, filter out aborted and not-yet committed transactions), to consume all messages, including for aborted transactions, you may set this property to read_uncommitted to get the behaviour of previous releases.
For consumers in read_committed mode, the end of a partition is now defined to be the offset of the last message of a successfully committed transaction (referred to as the 'Last Stable Offset').
For non-transactional messages there is no change from previous releases, they will always be read, but a consumer will not advance into a not yet committed transaction on the partition.
linger.ms default was changed from 0 to 0.5 ms to promote some level of batching even with default settings.isolation.level=read_committed ensures the consumer will only read messages from successfully committed producer transactions. Default is read_committed. To get the previous behaviour, set the property to read_uncommitted, which will read all messages produced to a topic, regardless if the message was part of an aborted or not yet committed transaction.General:
linger.ms, this reduces CPU load and lock contention for high throughput producer applications. (#2509)enable.ssl.certificate.verification=false (@salisbury-espinosa)Consumer:
pause|resume() synchronous, ensuring that a subsequent poll() will not return messages for the paused partitions.Producer:
message.timeout.ms=0 is now accepted even if linger.ms > 0 (by Jeff Snyder)confluent-kafka-python is based on librdkafka v1.1.0, see the librdkafka v1.1.0 release notes for a complete list of changes, enhancements, fixes and
confluent-kafka-python is based on librdkafka v1.1.0, see the librdkafka v1.1.0 release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
ssl.endpoint.identification.algorithm=https (off by default) to validate the broker hostname matches the certificate. Requires OpenSSL >= 1.0.2(included with Wheel installations))ssl.ca.location), librdkafka will load the CA certs by default from the Windows Root Certificate Store.enable.ssl.certificate.verification=false)%{broker.name} is no longer supported in sasl.kerberos.kinit.cmd since kinit refresh is no longer executed per broker, but per client instance.New configuration properties:
ssl.key.pem - client's private key as a string in PEM formatssl.certificate.pem - client's public key as a string in PEM formatenable.ssl.certificate.verification - enable(default)/disable OpenSSL's builtin broker certificate verification.enable.ssl.endpoint.identification.algorithm - to verify the broker's hostname with its certificate (disabled by default).rd_kafka_conf_set_ssl_cert() to pass PKCS#12, DER or PEM certs in (binary) memory form to the configuration object.message.timeout.ms max value from 15 minutes to 24 days (@sarkanyi, workaround for #2015)sasl.kerberos.kinit.cmd to first attempt ticket refresh, then acquire.max.poll.interval.ms now correctly handles blocking poll calls, allowing a longer poll timeout than the max poll interval.confluent-kafka-python is based on librdkafka v1.0.1, see the librdkafka v1.0.1 release notes for a complete list of changes, enhancements, fixes and
confluent-kafka-python is based on librdkafka v1.0.1, see the librdkafka v1.0.1 release notes for a complete list of changes, enhancements, fixes and upgrade considerations.
v1.0.1 is a maintenance release with the following fixes:
This release also changes configuration defaults and deprecates a set of configuration properties, make sure to read the Upgrade considerations sectio…
confluent-kafka-python is based on librdkafka v1.0.0, see the librdkafka v1.0.0 release notes for a complete list of changes, enhancements and fixes and upgrade considerations.
v1.0.0 is a major feature release:
max.poll.interval.ms support in the Consumer.This release also changes configuration defaults and deprecates a set of configuration properties, make sure to read the Upgrade considerations section below.
The following configuration properties have changed default values, which may require application changes:
acks(alias request.required.acks) now defaults to all; wait for all in-sync replica brokers to ack. The previous default, 1 , only waited for an ack from the partition leader. This change places a greater emphasis on durability at a slight cost to latency. It is not recommended that you lower this value unless latency takes a higher precedence than data durability in your application.
broker.version.fallback now to defaults to 0.10, previously 0.9. broker.version.fallback.ms now defaults to 0. Users on Apache Kafka <0.10 must set api.version.request=false and broker.version.fallback=.. to their broker version. For users >=0.10 there is no longer any need to specify any of these properties.
enable.partition.eof now defaults to false. KafkaError._PARTITION_EOF was previously emitted by default to signify the consumer has reached the end of a partition. Applications which rely on this behavior must now explicitly set enable.partition.eof=true if this behavior is required. This change simplifies the more common case where consumer applications consume in an endless loop.
group.id is now required for Python consumers.
The following configuration properties have been deprecated. Use of any deprecated configuration property will result in a warning when the client instance is created. The deprecated configuration properties will be removed in a future release.
offset.store.method=file is deprecated.offset.store.path is deprecated.offset.store.sync.interval.ms is deprecated.produce.offset.report is no longer used. Offsets are always reported.queuing.strategy was an experimental property that is now deprecated.reconnect.backoff.jitter.ms is no longer used, see reconnect.backoff.ms and reconnect.backoff.max.ms.socket.blocking.max.ms is no longer used.topic.metadata.refresh.fast.cnt is no longer used.default.topic.config is deprecated.This release adds support for Idempotent Producer, providing exactly-once producing and guaranteed ordering of messages.
Enabling idempotence is as simple as setting the enable.idempotence
configuration property to true.
There are no required application changes, but it is recommended to add support for the newly introduced fatal errors that will be triggered when the idempotent producer encounters an unrecoverable error that would break the ordering or duplication guarantees.
See Idempotent Producer in the manual and the Exactly once semantics blog post for more information.
In previous releases librdkafka would maintain open connections to all brokers in the cluster and the bootstrap servers.
With this release librdkafka now connects to a single bootstrap server to retrieve the full broker list, and then connects to the brokers it needs to communicate with: partition leaders, group coordinators, etc.
For large scale deployments this greatly reduces the number of connections between clients and brokers, and avoids the repeated idle connection closes for unused connections.
Sparse connections is on by default (recommended setting), the old
behavior of connecting to all brokers in the cluster can be re-enabled
by setting enable.sparse.connections=false.
See Sparse connections in the manual for more information.
Original issue librdkafka #825.
max.poll.interval.ms is enforcedThis release adds support for max.poll.interval.ms (KIP-62), which requires
the application to call consumer.poll() at least every max.poll.interval.ms.
Failure to do so will make the consumer automatically leave the group, causing a group rebalance,
and not rejoin the group until the application has called ..poll() again, triggering yet another group rebalance.
max.poll.interval.ms is set to 5 minutes by default.
CachedSchemaRegistryClient configuration with configuration dict for application configsDelete Schema support to CachedSchemaRegistryClientConsumer.consume without setting group.id(now required)CachedSchemaRegistryClient handles get_compatibility properly./tests/run.sh added to simplify unit and integration test executionNothing published for this version
The property default.topic.configuration has been deprecated and will be removed in 1.0, but still has precedence to topic configuration specified in…
See librdkafka v0.11.6 release notes for enhancements and fixes in librdkafka.
default.topic.configuration has been deprecated and will be removed in 1.0, but still has precedence to topic configuration specified in the global configuration dictionary. (#446)debug configuration property prior to plugin.library.paths for enhanced debugging. (#464)Nothing published for this version
v0.11.5 is a feature release that adds support for the Kafka Admin API (KIP-4).
v0.11.5 is a feature release that adds support for the Kafka Admin API (KIP-4).
This release adds support for the Admin API, enabling applications and users to perform administrative Kafka tasks programmatically:
The API closely follows the Java Admin API:
def example_create_topics(a, topics):
new_topics = [NewTopic(topic, num_partitions=3, replication_factor=1) for topic in topics]
# Call create_topics to asynchronously create topics
fs = a.create_topics(new_topics)
# Wait for operation to finish.
for topic, f in fs.items():
try:
f.result() # The result itself is None
print("Topic {} created".format(topic))
except Exception as e:
print("Failed to create topic {}: {}".format(topic, e))
Additional examples can be found in examples/adminapi
throttle_cb (#237) (#377)test_compatibility() should return False not None would return None when unable to check compatibility (#372, @Enether)Producer.produce documentation to use correct time unit of seconds (#384) (#385)Your coding agent can read these notes before it upgrades. Set up the MCP server →