v2.16.0 is a feature release with the following features, fixes and enhancements:
Enhancements
confluent_kafkanow 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)- Add Python 3.14t wheels (#2352)
- Add Linux s390x (IBM Z) wheels, including Python 3.14t (#2360)
- Producer
close()now aborts any open transaction (#2347) - Async IO Consumer's default worker pool size has been increased from 2 to 100 (#2347)
SerializingProducerandDeserializingConsumeraccept a serde builder
when not passing a ready-made serde, through the newkey.serializer.builder/
value.serializer.builderandkey.deserializer.builder/
value.deserializer.builderconfiguration properties. Serdes built this way,
and any Schema Registry client the builder created for them, are owned by the
client and closed by itsclose(); ready-made serdes remain the
application's. This API is experimental (#2364).- Serdes that resolve subjects through the Schema Registry associated subject
name strategy are now given the Kafka cluster id automatically. The id is
resolved lazily, on the first subject lookup, so creating a client never
waits on a broker, and concurrent lookups share a singlecluster_id()
call; until a broker has been reached the lookup raises a
SerializationErrornamingsubject.name.strategy.kafka.cluster.id, which
can be set to supply the id explicitly. The new serde hooks
(set_cluster_id_resolver(),close()) are experimental (#2364). - New
Producer.cluster_id(),Consumer.cluster_id()and
AdminClient.cluster_id()(also onAIOProducerandAIOConsumer),
returning the id of the cluster the client is connected to. This API is
experimental (#2364). - New asyncio clients
AsyncSerializingProducerandAsyncDeserializingConsumer
inconfluent_kafka.aio, the counterparts ofSerializingProducerand
DeserializingConsumerbuilt onAIOProducer/AIOConsumer, accepting the
asyncio Schema Registry serdes and their builders. These classes are
experimental (#2364). - New
Message.deserialized_key()andMessage.deserialized_value(), which
return the same objects askey()andvalue()but are typed with the
deserialized types on aDeserializingConsumer.SerializingProducerand
DeserializingConsumerare now generic in their key and value types. This
API is experimental (#2364). - Add support for saving Azure key version with DEK (#2306)
- Pass context when clients make KEK calls to DEK Registry (#2308)
- Schema Registry: add support for the DLQ (dead-letter-queue) rule action
(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)RuleRegistrythe DLQ is best-effort; set
dlq.auto.flush=trueor give the serde its ownRuleRegistry(closable on
shutdown) for durability. - Add support for inline validation rules (#2326)
- Add Variant, Decimal, and Timestamp CEL functions (#2332)
Fixes
- Raise
SerializationErrorfor Unicode encoding and decoding failures in string serializers. - Fix concurrency safety issues in
Producer,Consumer,AdminClient, and
Messageclasses (#2347) - Serialize concurrent access to a shared
Consumerinstance across threads
instead of leaving it as undefined behavior (#2347) - Prefer httpx2 over httpx for Schema Registry to avoid Authlib deprecation warnings (#2351)
- Fix KafkaError error strings raising/garbling on non-UTF-8 locales (#2331)
- Fix crash on nullable array of $ref items in JSON Schema CSFLE (#2370)
- Fix
Producer.purge()ignoringin_queue,in_flightandblockingset toFalseon big-endian platforms such as s390x (#2345)
confluent-kafka-python 2.16.0 is based on librdkafka 2.16.0, see the librdkafka release notes for a complete list of changes, enhancements, fixes and upgrade considerations.