NewYour coding agent can read the release notes before it upgrades.Set up the MCP server →
PyPI · #859 most downloaded on PyPI
Pure Python client for Apache Kafka
Last release 1 months ago
16 Aug 2026
Ships unpredictably
gaps range from 8 days to 5 months
Nearly every release is documented
notes for 60 of the last 60 stable releases
Nothing withdrawn
no release was ever pulled
13 years old
80 releases · first in 2014
Fix JsonSerializer returning None and raising TypeError on deserialize
Admin: Fix create_partitions sending empty assignments array instead of null
create_partitions sending empty assignments array instead of null (#3148)One column per quarter.
Consumer: fix stuck partitions retained after rebalance with listener
transport: Fix TLS SNI when ssl_check_hostname=False
transport: cancel pending io coroutines on close()
consumer: track current leader epoch in addition to record epoch
Reserve buffer capacity before every encode_into write
Fix _build_transport / conn.close race
_build_transport / conn.close race (#3097)_write_to_sock (#3095)fix memory leak of cancelled timeout tasks ( @azdobylak / #3077 )
fix memory leak of cancelled timeout tasks ([azdobylak](https://github.com/azdobylak) / #3077)
receive_message_max_bytes default 100MiB (consumer/producer/admin); check against fetch_max_bytes and max_partition_fetch_bytes in Consumer
It substantially refactors and expands the Admin client, including breaking changes to some API signatures, and it lands a long list of KIP features/c…
This is a major release with significant changes to kafka-python internals to simplify networking and feature development. It introduces a new networking layer (kafka.net) and a dynamic protocol system that uses JSON schema files imported from Apache Kafka. It substantially refactors and expands the Admin client, including breaking changes to some API signatures, and it lands a long list of KIP features/changes across the producer, consumer, admin, and networking/metadata clients. Protocol support across kafka-python is now at or beyond the apache kafka 3.0 baseline.
Complete refactor of the networking layer using a bespoke event-loop supporting async/await (but no asyncio yet). All three clients (Admin, Consumer, Producer) use a dedicated IO thread that drives a selector-based event loop. AdminClient/KafkaConsumer leverage a built-in io thread supplied by kafka.net, KafkaProducer continues to use its existing background Sender thread for now.
Per-request and per-stage timeouts replace the old single client-wide timeout.
The Future primitive gains __await__ and a faster slotted implementation; cross-thread wakeups are factored out into a reusable helper.
Defensive checks throughout the kafka.net event loop and transport stack: improved socket I/O error handling, RuntimeErrors on misuse of the IO thread, and lock-based detection of concurrent access.
_poll_once; add net.drain() (#2949)A new JSON-schema-based dynamic protocol generator now replaces the legacy hand-written protocol classes (moved to kafka.protocol.old).
Protocol classes are now generated from the upstream Apache Kafka JSON schemas.
Broker version inference is consolidated into a single BrokerVersionData helper that tracks the broker's reported API versions and infers a broker version string. ApiVersionsRequest is always sent on connect.
KafkaConsumer drops the dedicated HeartbeatThread in favor of scheduled async tasks on the kafka.net IO thread. Internals have been substantially refactored to migrate from future callbacks to async/await syntax. Feature support added for incremental cooperative rebalancing (KIP-429), rack-aware fetch (KIP-392), and log truncation detection (KIP-320/KIP-595).
All consumer network I/O now flows through the shared kafka.net IO thread; consumer.poll() no longer drives the event loop directly.
_update_fetch_positions -> _refresh_committed_offsets; dont poll in position() (#2958)Fetcher._fetch_offsets_by_times retry handling (#2833)KafkaProducer gains a sticky partitioner (KIP-480), enabled-by-default idempotence (KIP-679), tightened transaction handling, and a faster send/encode path.
Split-and-resend oversized batches instead of failing; avoid redundant validation and buffer copies on hot send-path.
Split KafkaAdminClient into focused mixin classes (cluster, topics, configs, groups, ACLs, log dirs, etc), and convert request-sending path to async def methods that run on the kafka.net IO thread. Support for new KIPs using new protocol stack.
The admin client interface remains sync but wraps a fully-async internal api (does not support asyncio yet). Adds cached coordinator lookups and a mixin structure to separate logical resource groups.
_send_request_to_controller error handling (#2751)The CLI adds shared parser config, SASL/SSL connection support across all subcommands, and several new admin subcommands (acls, configs alter, users).
Small quality-of-life additions to the public API surface.
Codec and Python-3-compatibility fixes that aren't specific to a single client.
A new in-memory MockBroker / MockTransport enables deterministic protocol-level tests, and the integration test fixtures have been substantially consolidated.
kafka.conn: Improve error handling for sasl authenticate mechanisms
_callbacks/_errbacks list when future is_done to avoid reference cycles (#2891)Fix TaggedFields value encoding; add test coverage
Fetcher._fetch_offsets_by_times retry handling (#2833)*2.3.x will be the last release branch with python2 support!*
2.3.x will be the last release branch with python2 support!
send_request() and send_requests() to KafkaAdminClient (#2649)kafka.conn: Improve error handling for sasl authenticate mechanisms
_callbacks/_errbacks list when future is_done to avoid reference cycles (#2891)Fix TaggedFields value encoding; add test coverage
Fetcher._fetch_offsets_by_times retry handling (#2833)Add ProducerBatch.__lt__ for heapq
Fixes
Add internal poll to consumer.position()
Fixes
Networking
Documentation
transactional_id to KafkaProducer Keyword Arguments docstringFix thread not waking up when there is still data to be sent (gqmelo / #2670)
Fixes
Consumer.position() (k61n / #2668)Fix KafkaProducer broken method names (@llk89 / #2660)
Fixes
Fix coordinator lock contention during close()
Use client.await_ready() to simplify blocking wait and add timeout to admin client
Fix construction of final GSSAPI authentication message
_completed_fetches deque in consumer fetcher (#2646)Do not ignore metadata response for single topic with error
Set the current host in the SASL configs
client_ctx.complete in auth_bytes() (#2631)Do not reset fetch positions if offset commit fetch times out
Wait for next heartbeat in thread loop; check for connected coordinator
Minor Heartbeat updates: catch more exceptions / log configuration / raise KafkaConfigurationError
Only disable heartbeat thread once at beginning of join-group
Fix producer busy loop with no pending batches
Do not reset_generation after RebalanceInProgressError; improve CommitFailed error messages
reset_generation after RebalanceInProgressError; improve CommitFailed error messages (#2614)Ignore leading SECURITY_PROTOCOL:// in bootstrap_servers
# 2.2.2 (Apr 30, 2025) Fixes ----- * Fix lint errors
Always try ApiVersionsRequest v0, even on broker disconnect
Potentially Breaking Changes (internal)
delivery_timeout_ms_tp_locks; get dq with lock in reenqueue()READ_COMMITTED (#2582)MEMBER_ID_REQUIRED error w/ second join group request (#2598)ClusterMetadata.add_group_coordinator -> add_coordinator + support txn typelog_start_offset from producer RecordMetadataDefaultRecordsBuilder.size_in_bytes to classmethodtest_fetcherOnly create fetch requests for ready nodes
Move benchmark scripts to kafka.benchmarks module
__slots__ for metrics (#2583)metrics_enabled=False to disable metrics (#2581)Try import new Sequence before old to avoid DeprecationWarning
Fix crash when switching to closest compatible api_version in KafkaClient
Simplify consumer.poll send fetches logic
_unpack_records in PartitionRecords to fix premature fetch offset advance in consumer.poll() (#2555)Fix packaging of 2.1.0 in Fedora: testing requires "pytest-timeout".
Support Kafka Broker 2.1 API Baseline
_maybe_auto_commit_offsets_async (#2546)Improve error handling in client._maybe_connect
client._maybe_connect (#2504)maybe_refresh_metadata changes (#2507)client.check_version timeout to api_version_auto_timeout_ms (#2496)__del__Remove unused client bootstrap backoff code
servers/*/api_versions and servers/*/messagesCheck for wakeup socket errors on read and close and reinit to reset
Update socketpair w/ CVE-2024-3219 fix
Fix crc32c deprecation warning (crc32c==2.1) (jeffwidman / PR #2128)
KAFKA-8962: Use least_loaded_node() for AdminClient.describe_topics() (jeffwidman / PR #2000)
This release includes breaking changes for any application code that has not migrated from older Simple-style classes to newer Kafka-style classes.
This release includes breaking changes for any application code that has not migrated from older Simple-style classes to newer Kafka-style classes.
IncompatibleBrokerVersion when passing an api_version (ian28223 / PR #1953)ConnectionError (jeffwidman / PR #1816)This release is focused on KafkaConsumer performance, Admin Client improvements, and Client concurrency. The KafkaConsumer iterator implementation has
This release is focused on KafkaConsumer performance, Admin Client
improvements, and Client concurrency. The KafkaConsumer iterator implementation
has been greatly simplified so that it just wraps consumer.poll(). The prior
implementation will remain available for a few more releases using the optional
KafkaConsumer config: legacy_iterator=True . This is expected to improve
consumer throughput substantially and help reduce heartbeat failures / group
rebalancing.
Major thanks to @carsonip @Baisang @iv-m @davidheitman @cardy31 @ulrikjohansson @iAnomaly @Wayde2014 @ossdev07 @commanderdishwasher @justecorruptio @melor @rustyrothwurt @sachiin @jacky15 and @rikonen for submitting PRs; thanks as well to everyone that submitted bug reports and issues, and to @jeffwidman and @tvoinarovskyi for code reviews, comments, testing, debugging, and helping to maintain kafka-python!
consumer.poll() for KafkaConsumer iteration (@dpkp / PR #1902)partitions_for_topic a read-through cache (@Baisang / PR #1781,#1809)ssl_cafile is not provided (@iAnomaly / PR #1883)sasl_kerberos_domain_name config to KafkaAdminClient (@jeffwidman / PR #1852)security_protocol config documentation for KafkaAdminClient (@cardy31 / PR #1849)_send_request_to_node() in KafkaAdminClient (@davidheitman / PR #1807)KafkaConsumer tests to pytest (@jeffwidman / PR #1886)KAFKA_VERSION env var in tests (@jeffwidman / PR #1887)socket.SOCK_STREAM in test assertions (@iv-m / PR #1879)consumer.topics() and consumer.partitions_for_topic() (@Baisang / PR #1829)api_version_auto_timeout_ms (@jeffwidman / PR #1812)This is a minor release focused on KafkaConsumer performance, Admin Client
improvements, and Client concurrency. The KafkaConsumer iterator implementation
has been greatly simplified so that it just wraps consumer.poll(). The prior
implementation will remain available for a few more releases using the optional
KafkaConsumer config: legacy_iterator=True . This is expected to improve
consumer throughput substantially and help reduce heartbeat failures / group
rebalancing.
Client
KafkaConsumer
consumer.poll() for KafkaConsumer iteration (dpkp / PR #1902)partitions_for_topic a read-through cache (Baisang / PR #1781,#1809)Miscellaneous Bugfixes / Improvements
ssl_cafile is not provided (iAnomaly / PR #1883)Admin Client
sasl_kerberos_domain_name config to KafkaAdminClient (jeffwidman / PR #1852)security_protocol config documentation for KafkaAdminClient (cardy31 / PR #1849)_send_request_to_node() in KafkaAdminClient (davidheitman / PR #1807)Test Infrastructure / Documentation / Maintenance
KafkaConsumer tests to pytest (jeffwidman / PR #1886)KAFKA_VERSION env var in tests (jeffwidman / PR #1887)socket.SOCK_STREAM in test assertions (iv-m / PR #1879)consumer.topics() and consumer.partitions_for_topic() (Baisang / PR #1829)api_version_auto_timeout_ms (jeffwidman / PR #1812)This is a patch release primarily focused on bugs related to concurrency, SSL connections and testing, and SASL authentication. Major thanks to @pt2ph
This is a patch release primarily focused on bugs related to concurrency, SSL connections and testing, and SASL authentication. Major thanks to @pt2pham , @isamaru , @braedon , @gingercookiemage , for submitting PRs to help fix many of these issues. And major thanks to everyone that submitted bug reports and issues. And thanks always to @jeffwidman and @tvoinarovskyi for code reviews, comments, testing, debugging, and helping to maintain this project!
protocol.send_bytes (isamaru / PR #1752)state_change_callback with lock (dpkp / PR #1775)client._conns in send() (dpkp / PR #1772)client.check_version (dpkp / PR #1771)maybe_refresh_metadata -- it is only called by poll() (dpkp / PR #1769)Travis CI: 'sudo' tag is now deprecated in Travis (cclauss / PR #1698)
This release is primarily focused on addressing lock contention and other coordination issues between the KafkaConsumer and the background heartbeat thread that was introduced in the 1.4 release.
skip_double_compressed_messages (jeffwidman / PR #1677)Stop using deprecated log.warn() (jeffwidman #1615)
validate_only/include_synonyms (jeffwidman #1645)kafka.common internally (jeffwidman #1509)pylint (jeffwidman #1611)Unittest to pytest (jeffwidman #1620)six to 1.11.0 (jeffwidman #1602)six consistently (jeffwidman #1605)pylint import errors on six.moves (jeffwidman #1609)Stop using deprecated log.warn() (jeffwidman)
ConnectionError (jeffwidman #1492)Close leaked selector in version check (dpkp #1425)
BrokerConnection.connection_delay() to return milliseconds (dpkp #1414)Fetcher._fetchable_partitions to avoid mutation errors (dpkp #1400)_unpack (j2gg0s #1403)BrokerConnection.connect_blocking() to improve bootstrap to multi-address hostnames (dpkp #1411)BrokerConnection.close() if already disconnected (dpkp #1424)api_version against known versions (dpkp #1434)max_records in KafkaConsumer.poll (dpkp #1398)Fix consumer poll stuck error when no available partition (ckyoog #1375)
Use non-deprecated exception handling (jeffwidman a699f6a)
This is a substantial release. Although there are no known 'showstopper' bugs as of release, we do recommend you test any planned upgrade to your application prior to running in production.
Some of the major changes include:
Thanks to all contributors -- the state of the kafka-python community is strong!
Detailed changelog are listed below:
Fix partition assignment race condition (jeffwidman #1240)
Avoid multiple connection attempts when refreshing metadata (dpkp #1067)
Derive all api classes from Request / Response base classes (dpkp 1030)
Core / Protocol
max_bytes option and FetchRequest_v3 usage. (Drizzt1991 962)Test Infrastructure
Consumer
Producer
Client
Bugfixes
Logging / Error Messages
Documentation
Legacy Client
Your coding agent can read these notes before it upgrades. Set up the MCP server →