From 3cc5a80944e6cfddf88108f8e4201600ff01fc40 Mon Sep 17 00:00:00 2001 From: Tommi Vainikainen Date: Fri, 15 May 2026 09:35:41 +0300 Subject: [PATCH] Bump kafka-python dependency to support 2.3.1 Update the version constraint from <2.0.4 to support latest 2.3.1 release. kafka-python 2.2.0+ removed errors.RETRY_ERROR_TYPES in favor of per-class retriable attributes. Add _get_retriable_kafka_errors() helper that falls back to introspecting error classes when the legacy constant is unavailable, maintaining compatibility with both old and new versions. Co-authored-by: Copilot <223556219+Copilot@users.noreply.github.com> --- journalpump/senders/kafka_sender.py | 19 ++++++++++++++++++- requirements.txt | 2 +- setup.py | 10 ++++++---- 3 files changed, 25 insertions(+), 6 deletions(-) diff --git a/journalpump/senders/kafka_sender.py b/journalpump/senders/kafka_sender.py index 51a3496..ac6cfb6 100644 --- a/journalpump/senders/kafka_sender.py +++ b/journalpump/senders/kafka_sender.py @@ -2,6 +2,7 @@ from kafka import errors, KafkaAdminClient, KafkaProducer from kafka.admin import NewTopic +import inspect import logging import socket @@ -15,7 +16,23 @@ except ImportError: zstd = None -KAFKA_CONN_ERRORS = tuple(errors.RETRY_ERROR_TYPES) + ( + +def _get_retriable_kafka_errors(): + """Build the set of retriable Kafka error types. + + kafka-python < 2.1 provides errors.RETRY_ERROR_TYPES. + kafka-python >= 2.1 marks individual error classes with retriable = True. + """ + if hasattr(errors, "RETRY_ERROR_TYPES"): + return tuple(errors.RETRY_ERROR_TYPES) + return tuple( + obj + for _name, obj in inspect.getmembers(errors) + if inspect.isclass(obj) and issubclass(obj, errors.KafkaError) and getattr(obj, "retriable", False) + ) + + +KAFKA_CONN_ERRORS = _get_retriable_kafka_errors() + ( errors.UnknownError, socket.timeout, TimeoutError, diff --git a/requirements.txt b/requirements.txt index 3f9fbe7..c52d43f 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1,4 +1,4 @@ -kafka-python<2.0.4 +kafka-python<2.4.0 requests websockets<14.0 aiohttp-socks diff --git a/setup.py b/setup.py index 9353d30..b076430 100644 --- a/setup.py +++ b/setup.py @@ -1,7 +1,8 @@ -from journalpump import __version__ +import os + from setuptools import find_packages, setup -import os +from journalpump import __version__ setup( name="journalpump", @@ -10,10 +11,11 @@ packages=find_packages(exclude=["test", "systest"]), extras_require={}, install_requires=[ - "kafka-python", + "kafka-python<2.4.0", "requests", - "websockets", + "websockets<14.0", "aiohttp-socks", + "python-snappy", "botocore", "google-api-python-client", "google-auth",