Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
19 changes: 18 additions & 1 deletion journalpump/senders/kafka_sender.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
from kafka import errors, KafkaAdminClient, KafkaProducer
from kafka.admin import NewTopic

import inspect
import logging
import socket

Expand All @@ -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,
Expand Down
2 changes: 1 addition & 1 deletion requirements.txt
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
kafka-python<2.0.4
kafka-python<2.4.0
requests
websockets<14.0
aiohttp-socks
Expand Down
10 changes: 6 additions & 4 deletions setup.py
Original file line number Diff line number Diff line change
@@ -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",
Expand All @@ -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",
Expand Down
Loading