BI Consult Desktop Logo BI Consult Mobile Logo
  • Russian BI Исследование российских bi
  • Перейти на Fine BI
  • Контакты
  • +7 812 334-08-01
    +7 499 608-13-06
  • Отправить сообщение
  • Главная
  • Продукты Эксперт-BI
    • Дистрибуция
    • Розничная торговля
    • Производство
    • Операторы связи
    • Страхование
    • Банки
    • Лизинг
    • Логистика
    • Нефтегазовый сектор
    • Медицина
    • Сеть ресторанов
    • E-Commerce
    • Сельское хозяйство
    • Энергетика
    • FMCG
    • Девелоперы
    • Маркетплейсы
    • Пищевая промышленность
    • Фармацевтика
    • Построение Data Platform
    • Цифровая трансформация
    • Управление по KPI
    • Финансы
    • Продажи
    • Склад
    • HR
    • Маркетинг
    • Внутренний аудит
    • Категорийный менеджмент
    • S&OP и FP&A
    • Геоаналитика
    • Цепочки поставок (SCM)
    • AutoML
    • Process Mining
    • IBP
    • ИТ (CIO)
    • Закупки
  • Платформы
    • Системы бизнес-анализа (BI)
    • Интегрированное бизнес-планирование (IBP)
    • Хранилища данных (DWH / Lakehouse)
    • Каталоги данных (Data Catalog)
    • Системы ETL и ELT
    • AI / Исскуственный интеллект
    • Шина данных (ESB)
    • Система управления мастер-данными (MDM)
    • Семантический слой
  • Услуги
    • Переход на отечественные BI и DWH системы
    • Консалтинг
    • Пилотный проект
    • Обучение и сертификация
    • Бесплатное обучение
    • Поддержка
    • Технические задания
    • Сбор требований для проекта внедрения BI-системы
    • CI/CD для DWH
    • Аудит BI приложений и DWH
    • Выделенная команда
    • Настойка и поддержка баз данных
    • Разработка BI Стратегии
    • Styleguide для BI-системы
    • Как выбрать BI-систему
  • Курсы
    • Учебный курс Информационная грамотность (Data Literacy)
    • Учебный курс для бизнес-аналитиков
    • Учебный курс для системных аналитиков
    • Учебный курс по Data Governance
    • Учебный курс Как стать CDO
    • Учебный курс Современная архитектура хранилища данных
    • Учебный курс по Fine BI
    • Учебный курс по FineReport
    • Учебный курс по DWH
    • Учебный курс по Data Science (ML, AI)
    • Учебный курс по PostgreSQL
    • Учебный курс по Greenplum
    • Учебный курс по Apache Airflow и NiFi
    • Учебный курс по Open-source BI
    • Учебный курс по ClickHouse
    • Учебный курс по DataLens
    • Учебный курс по Loginom
    • Учебный курс по Modus BI и ETL
    • Учебный курс по Visiology
    • Учебный курс по dbt (Data Build Tool)
  • Компания
    • Руководство
    • Новости
    • Клиенты
    • Карьера
    • Скачать
    • Контакты

BI

  • FineBI
  • FineReport
  • FineDataLink
  • FineChatBI (FineAI)
  • Коннекторы данных из 1С в BI
  • Airflow / Nifi
  • Visiology
  • PIX BI
  • Modus BI
  • Yandex.DataLens
  • Open-source BI: Superset/Metabase
  • Luxms BI
  • AW BI + Alpha BI
  • FlyBI + Форсайт. Аналитическая Платформа
  • Loginom
  • Триафлай
  • AI / Исскуственный интеллект
  • Optimacros
  • Навигатор BI
  • Семантический слой

СУБД

  • Arenadata
  • ClickHouse
  • Greenplum
  • Postgres Professional
  • TData

Другое

  • Построение Data Platform
    • Аналитическое хранилище данных
    • Data Lake и Data Engineering
    • Подробнее про Data Lake
    • Внедрение Lakehouse
      • Apache Doris
      • StarRocks
      • Trino
    • Миграция витрин из пропиетарных DWH на новый стек
    • Учебный курс "Современная архитектура хранилища данных"
Главная » Курсы по системам бизнес-анализа и методологии » Учебный курс Современная архитектура хранилища данных » kafka-python

kafka-python

Введение

kafka-python разработан таким образом, чтобы быть в состоянии функционировать подобно официальному java-клиенту, только с добавлением питоновских интерфейсов (например, потребительских итераторов).

kafka-python лучше всего использовать с новыми брокерами (0.9+), но он совместим и с более старыми версиями (до 0.8.0). Некоторые функции будут доступны только в новых брокерах. Например, полностью скоординированные группы потребителей - т.е. динамическое назначение разделов нескольким потребителям в одной группе - требуют использования брокеров kafka версии 0.9. Применение этой функции в более ранних версиях брокера потребует написания специального кода (возможно, с использованием zookeeper или consul). Для более старых брокеров Вы можете добиться чего-то подобного, вручную назначив различные разделы каждому экземпляру потребителя с помощью инструментов управления конфигурацией, таких как chef, ansible и т. д. Такой подход сработает несмотря на то, что он не поддерживает ребалансировку при сбоях. Более подробную информацию см. в разделе «Совместимость».

Обратите внимание, что мастер-ветка может содержать невыпущенные функции. Документацию по релизу можно найти в readthedocs и/или встроенной справке python.

>>> pip inst​all kafka-python

 

KafkaConsumer

KafkaConsumer - это высокоуровневый консюмер сообщений. Полная поддержка скоординированных групп консюмеров требует использования брокеров kafka, поддерживающих Group API: kafka v0.9+.

Итератор консюмеров возвращает записи ConsumerRecords, которые представляют собой простые кортежи, раскрывающие основные атрибуты сообщений: тему, раздел, смещение, ключ и значение:

>>> from kafka import KafkaConsumer
>>> consumer = KafkaConsumer('my_favorite_topic')
>>> for msg in consumer:
...     print (msg)
>>> # join a consumer group for dynamic partition assignment and offset commits
>>> from kafka import KafkaConsumer
>>> consumer = KafkaConsumer('my_favorite_topic', group_id='my_favorite_group')
>>> for msg in consumer:
...     print (msg)
>>> # manually assign the partition list for the consumer
>>> from kafka import TopicPartition
>>> consumer = KafkaConsumer(bootstrap_servers='localhost:1234')
>>> consumer.assign([TopicPartition('foobar', 2)])
>>> msg = next(consumer)
>>> # Deserialize msgpack-encoded values
>>> consumer = KafkaConsumer(value_deserializer=msgpack.loads)
>>> consumer.subscribe(['msgpackfoo'])
>>> for msg in consumer:
...     assert isinstance(msg.value, dict)

 

KafkaProducer

KafkaProducer - это высокоуровневый асинхронный продюсер сообщений.

>>> from kafka import KafkaProducer
>>> producer = KafkaProducer(bootstrap_servers='localhost:1234')
>>> for _ in range(100):
...     producer.send('foobar', b'some_message_bytes')
>>> # Block until a single message is sent (or timeout)
>>> future = producer.send('foobar', b'another_message')
>>> result = future.get(timeout=60)
>>> # Block until all pending messages are at least put on the network
>>> # NOTE: This does not guarantee delivery or success! It is really
>>> # only useful if you configure internal batching using linger_ms
>>> producer.flush()
>>> # Use a key for hashed-partitioning
>>> producer.send('foobar', key=b'foo', value=b'bar')
>>> # Serialize json messages
>>> import json
>>> producer = KafkaProducer(value_serializer=lambda v: json.dumps(v).encode('utf-8'))
>>> producer.send('fizzbuzz', {'foo': 'bar'})
>>> # Serialize string keys
>>> producer = KafkaProducer(key_serializer=str.encode)
>>> producer.send('flipflap', key='ping', value=b'1234')
>>> # Compress messages
>>> producer = KafkaProducer(compression_type='gzip')
>>> for i in range(1000):
...     producer.send('foobar', b'msg %d' % i)

 

Использование

KafkaProducer можно использовать в разных потоках (в отличие от KafkaConsumer).

Хотя KafkaConsumer можно использовать в локальном потоке, все же лучше использовать многопоточную обработку.

 

Сжатие

kafka-python поддерживает gzip-сжатие/декомпрессию. Чтобы создавать или потреблять сообщения, сжатые по протоколу lz4, необходимо установить python-lz4 (pip install lz4). Чтобы включить snappy, установите python-snappy (также требуется библиотека snappy). 

 

Протокол

Вторичная цель kafka-python - предоставить простой в использовании протокольный слой для взаимодействия с брокерами kafka через python repl. Это полезно для тестирования, зондирования и общих экспериментов. Поддержка протокола используется для включения метода check_version(), который проверяет брокер kafka и пытается определить, какую именно  версию он использует (от 0.8.0 до 2.4+).

 

Обзор KafkaConsumer

from kafka import KafkaConsumer
 
# To consume latest messages and auto-commit offsets
consumer = KafkaConsumer('my-topic',
                         group_id='my-group',
                         bootstrap_servers=['localhost:9092'])
for message in consumer:
    # message value and key are raw bytes -- decode if necessary!
    # e.g., for unicode: `message.value.decode('utf-8')`
    print ("%s:%d:%d: key=%s value=%s" % (message.topic, message.partition,
                                          message.offset, message.key,
                                          message.value))
 
# consume earliest available messages, don't commit offsets
KafkaConsumer(auto_offset_reset='earliest', enable_auto_commit=False)
 
# consume json messages
KafkaConsumer(value_deserializer=lambda m: json.loads(m.decode('ascii')))
 
# consume msgpack
KafkaConsumer(value_deserializer=msgpack.unpackb)
 
# StopIteration if no message after 1sec
KafkaConsumer(consumer_timeout_ms=1000)
 
# Subscribe to a regex topic pattern
consumer = KafkaConsumer()
consumer.subscribe(pattern='^awesome.*')
 
# Use multiple consumers in parallel w/ 0.9 kafka brokers
# typically you would run each on a different server / process / CPU
consumer1 = KafkaConsumer('my-topic',
                          group_id='my-group',
                          bootstrap_servers='my.server.com')
consumer2 = KafkaConsumer('my-topic',
                          group_id='my-group',
                          bootstrap_servers='my.server.com')

 

Существует множество вариантов конфигурации класса консюмера. Более подробная информация приведена в документации по API KafkaConsumer.

 

KafkaProducer 

from kafka import KafkaProducer
from kafka.errors import KafkaError
 
producer = KafkaProducer(bootstrap_servers=['broker1:1234'])
 
# Asynchronous by default
future = producer.send('my-topic', b'raw_bytes')
 
# Block for 'synchronous' sends
try:
    record_metadata = future.get(timeout=10)
except KafkaError:
    # Decide what to do if produce request failed...
    log.exception()
    pass
 
# Successful result returns assigned partition and offset
print (record_metadata.topic)
print (record_metadata.partition)
print (record_metadata.offset)
 
# produce keyed messages to enable hashed partitioning
producer.send('my-topic', key=b'foo', value=b'bar')
 
# encode objects via msgpack
producer = KafkaProducer(value_serializer=msgpack.dumps)
producer.send('msgpack-topic', {'key': 'value'})
 
# produce json messages
producer = KafkaProducer(value_serializer=lambda m: json.dumps(m).encode('ascii'))
producer.send('json-topic', {'key': 'value'})
 
# produce asynchronously
for _ in range(100):
    producer.send('my-topic', b'msg')
 
def on_send_success(record_metadata):
    print(record_metadata.topic)
    print(record_metadata.partition)
    print(record_metadata.offset)
 
def on_send_error(excp):
    log.error('I am an errback', exc_info=excp)
    # handle exception
 
# produce asynchronously with callbacks
producer.send('my-topic', b'raw_bytes').add_callback(on_send_success).add_errback(on_send_error)
 
# block until all async messages are sent
producer.flush()
 
# configure multiple retries
producer = KafkaProducer(retries=5)

 

API KafkaConsumer

classkafka.KafkaConsumer(*topics, **configs)

Потребляйте записи из кластера Kafka.

Консюмер будет обрабатывать отказы серверов в кластере Kafka, а также адаптироваться к новым разделам топиков или в процессе миграции между брокерами. Он также взаимодействует с назначенным узлом Kafka Group Coordinator для того, чтобы позволить нескольким консюмерам сбалансировать нагрузку на потребление топиков (требуется kafka >= 0.9.0.0).

Конмюмер не должен быть общим для всех потоков.

Параметры

*topics (str) – необязательный список топиков для подписки. Если он не задан, вызовите subscribe() или assign() перед непосредственным потреблением записей.

 Ключевые аргументы

  • bootstrap_servers – ‘host[:port]’ string (или список строк ‘host[:port]’), к которому для загрузки начальных метаданных кластера должен обратиться консюмер. Это не обязательно должен быть полный список узлов. В нем должен быть хотя бы один брокер, который ответит на запрос Metadata API. По умолчанию используется порт 9092. Если серверы не указаны, по умолчанию будет использоваться локальный хост  9092.
  • client_id (str) – Имя для этого клиента. Эта строка передается в каждом запросе к серверам и может быть использована для идентификации конкретных записей журнала на стороне сервера, которые соответствуют этому клиенту. Также передается GroupCoordinator для ведения журнала в отношении администрирования групп потребителей. По умолчанию: ‘kafka-python-{version}’
  • group_id (str или None) – Имя группы консюмеров, к которой следует присоединиться для динамического назначения разделов (если включено) и использовать для получения и фиксации смещений. Если None, автоматическое назначение разделов (через координатора группы) и фиксация смещений отключены. По умолчанию: None
  • key_deserializer (вызываяемый) – Любой вызываемый модуль, который принимает необработанный ключ сообщения и возвращает десериализованный ключ.
  • value_deserializer (вызываемый – Any callable that takes a raw message value and returns a deserialized value.
  • fetch_min_bytes (int) – Минимальное количество данных, которое сервер должен вернуть по запросу fetch, в противном случае необходимо ждать, пока не накопится больше данных. По умолчанию: 1.
  • fetch_max_wait_ms (int) – Максимальное количество времени в миллисекундах, которое сервер будет блокировать перед ответом на запрос fetch, если данных недостаточно для немедленного удовлетворения требования, заданного параметром fetch_min_bytes. По умолчанию: 500.
  • fetch_max_bytes (int) – Максимальный объем данных, который сервер должен вернуть для запроса на выборку. Это не абсолютный максимум, если первое сообщение в первом непустом разделе выборки больше этого значения, сообщение все равно будет возвращено для того, чтобы консюмер мог продвигаться вперед. ПРИМЕЧАНИЕ: консюмер выполняет выборку для нескольких брокеров параллельно, поэтому использование памяти будет зависеть от количества брокеров, содержащих разделы для данного топика. Поддерживаемая версия Kafka >= 0.10.1.0. По умолчанию: 52428800 (50 МБ).
  • max_partition_fetch_bytes (int) – Максимальный объем данных для каждого раздела, который будет возвращен сервером. Максимальный общий объем памяти, используемый для запроса, = #partitions * max_partition_fetch_bytes. Этот размер должен быть не меньше максимального размера сообщения, разрешенного сервером, иначе консюмер будет отправлять слишком большие сообщения, которые не сможет получить консюмер. В этом случае консюмер может «застрять», пытаясь получить слишком большое для него сообщение на определенном разделе. По умолчанию: 1048576.
  • request_timeout_ms (int) – Таймаут запроса клиента в миллисекундах. По умолчанию: 305000.
  • retry_backoff_ms (int) – Миллисекунды для обратного хода при повторных попытках при ошибках. По умолчанию: 100.
  • reconnect_backoff_ms (int) – Количество времени в миллисекундах, которое необходимо выждать перед попыткой повторного подключения к данному узлу. По умолчанию: 50.
  • reconnect_backoff_max_ms (int) - Максимальное количество времени в миллисекундах для отката/ожидания при повторном подключении к брокеру, который неоднократно не смог установить соединение. Если указано, то время ожидания для каждого хоста будет экспоненциально увеличиваться для каждого последовательного отказа в соединении, вплоть до этого максимума. После достижения максимума попытки переподключения будут продолжаться периодически с этой фиксированной скоростью. Во избежание шторма соединений к коэффициенту обратного хода будет применен коэффициент рандомизации 0,2, в результате чего он будет находиться в случайном диапазоне между 20 % ниже и 20 % выше вычисленного значения. По умолчанию: 1000.
  • max_in_flight_requests_per_connection (int) – Запросы передаются брокерам kafka до указанного количества максимальных запросов на одно соединение с брокером. По умолчанию: 5.
  • auto_offset_reset (str) – Политика сброса смещений при ошибках OffsetOutOfRange: 'earliest' будет перемещаться к самому старому доступному сообщению, 'latest' - к самому последнему. Любое другое значение вызовет ошибку. По умолчанию: 'latest'.
  • enable_auto_commit (bool) – Если True, то смещение консюмера будет периодически фиксироваться в фоновом режиме. По умолчанию: True.
  • auto_commit_interval_ms (int) – Количество миллисекунд между автоматическими фиксациями смещения, если enable_auto_commit равно True. По умолчанию: 5000.
  • default_offset_commit_callback (вызываемый) – При вызове callback(offsets, response) ответом будет либо исключение, либо структура OffsetCommitResponse. Этот обратный вызов может быть использован для запуска пользовательских действий при завершении запроса фиксации.
  • check_crcs (bool) – Автоматическая проверка CRC32 потребляемых записей. Это гарантирует, что сообщения не были повреждены по проводам или на диске. Эта проверка увеличивает накладные расходы, поэтому ее можно отключить в случаях, требующих высокой производительности. По умолчанию: True
  • metadata_max_age_ms (int) – Период времени, исчисляемый в миллисекундах, по истечении которого мы принудительно обновляем метаданные, даже если мы не видели изменений в руководстве разделов, чтобы проактивно обнаружить новые брокеры или разделы. По умолчанию: 300000
  • partition_assignment_strategy (list) Список объектов, используемых для распределения владения разделами между экземплярами консюмеров при использовании группового управления. По умолчанию: [RangePartitionAssignor, RoundRobinPartitionAssignor].
  • max_poll_records (int) – Максимальное количество записей, возвращаемых за один вызов poll(). По умолчанию: 500
  • max_poll_interval_ms (int) – Максимальная задержка между вызовами poll() при использовании управления группами консюмеров. Устанавливает верхний предел времени, в течение которого консюмер может простаивать перед получением новых записей. Если poll() не будет вызван до истечения этого таймаута, то консюмер будет считаться вышедшим из строя, группа будет перебалансирована. По умолчанию 300000
  • session_timeout_ms (int) – Таймаут, используемый для обнаружения сбоев при использовании средств управления группами Kafka. Консюмер периодически посылает брокеру сигналы, сообщая тем самым о своей активности. Если до истечения этого таймаута брокер не получит ни одного сигнала, то брокер удалит этого консюмера из группы и инициирует перебалансировку. Обратите внимание на то, что значение данного параметра должно находиться в допустимом диапазоне, заданном в конфигурации брокера параметрами group.min.session.timeout.ms и group.max.session.timeout.ms. По умолчанию: 10000
  • heartbeat_interval_ms (int) – Ожидаемое время в миллисекундах между сигналами для координатора консюмеров при использовании средств управления группами Kafka. Сигналы используются для обеспечения того, чтобы сессия консюмера оставалась активной, а также для облегчения ребалансировки, когда новые консюмеры присоединяются или покидают группу. Это значение должно быть меньше, чем session_timeout_ms, но должно составлять не более 1/3 от этого значения. Оно может быть еще ниже для того, чтобы контролировать ожидаемое время ребалансировки. По умолчанию: 3000
  • receive_buffer_bytes (int) – Ожидаемое время в миллисекундах между ударами сердца для координатора потребителей при использовании средств управления группами Kafka. Сердцебиения используются для обеспечения того, чтобы сессия потребителя оставалась активной, и для облегчения ребалансировки, когда новые потребители присоединяются или покидают группу. Это значение должно быть меньше, чем session_timeout_ms, но обычно должно составлять не более 1/3 от этого значения. Его можно установить еще ниже, чтобы контролировать ожидаемое время нормальной ребалансировки. По умолчанию: 3000
  • send_buffer_bytes (int) – Размер буфера отправки TCP (SO_SNDBUF), который будет использоваться при отправке данных. По умолчанию: None (полагается на системные настройки по умолчанию). В java-клиенте по умолчанию используется значение 131072.
  • socket_options (list) – Список кортежей-аргументов для socket.setsockopt, применяемых к сокетам брокерских соединений. По умолчанию: [(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)].
  • consumer_timeout_ms (int) – количество миллисекунд, которое нужно блокировать во время итерации сообщения перед тем, как поднять StopIteration (т.е. завершить итератор). По умолчанию блокируется навсегда [float('inf')].
  • security_protocol (str) – Протокол, используемый для связи с брокерами. Допустимыми значениями являются: PLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL. По умолчанию: PLAINTEXT.
  • ssl_context (ssl.SSLContext) – Предварительно сконфигурированный SSLContext для обертывания сокетных соединений. Если указано, все остальные конфигурации ssl_* будут проигнорированы. По умолчанию: None.
  • ssl_check_hostname (bool) - Флаг для настройки того, должен ли ssl handshake проверять соответствие сертификата имени хоста брокера или нет. По умолчанию: True.
  • ssl_cafile (str) – Необязательное имя файла , используемого при проверке сертификата. По умолчанию: None.
  • ssl_certfile (str) – Необязательное имя файла в формате pem, содержащего сертификат клиента, а также любые сертификаты ca, необходимые для установления подлинности сертификата. По умолчанию: None.
  • ssl_keyfile (str) – Необязательное имя файла, содержащего закрытый ключ клиента. По умолчанию: None.
  • ssl_password (str) – Необязательный пароль, который будет использоваться при загрузке цепочки сертификатов. По умолчанию: None.
  • ssl_crlfile (str) – Необязательное имя файла, содержащего CRL для проверки срока действия сертификата. По умолчанию проверка CRL не производится. При указании файла только лист сертификата будет проверен по этому CRL. CRL может быть проверен только с Python 3.4+ или 2.7.9+. По умолчанию: None.
  • ssl_ciphers (str) –задает доступные шифры для ssl-соединений. Это должна быть строка в формате списка шифров OpenSSL. Если ни один шифр не может быть выбран (потому что опции времени компиляции или другая конфигурация запрещают использование всех указанных шифров), будет выдана ошибка ssl.SSLError. См. ssl.SSLContext.set_ciphers
  • api_version (tuple) – Укажите, какую версию API Kafka Вы собираетесь использовать. Если установлено значение None, клиент будет пытаться определить версию брокера путем опроса различных API. Разные версии обеспечивают различную функциональность.

Примеры

  • (0, 9) обеспечивает полную координацию действий группы с автоматическим партиционированием и ребалансировкой
  • (0, 8, 2) позволяет фиксировать смещение kafka-хранилища с назначением разделов вручную
  • (0, 8, 1) включает коммиты смещения zookeeper-storage с назначением разделов вручную
  • (0, 8, 0) обеспечивает базовую функциональность, но требует ручного назначение разделов и управление смещением.

 По умолчанию: None

  • api_version_auto_timeout_ms (int) – количество миллисекунд для выброса исключения таймаута из конструктора при проверке версии api брокера. Применяется, только если api_version установлено в None.
  • connections_max_idle_ms – Закрывает простаивающие соединения через количество миллисекунд, указанное в этом конфиге. Брокер закрывает простаивающие соединения после connections.max.idle.ms, что позволяет избежать неожиданных ошибок отключения сокетов на клиенте. По умолчанию: 540000
  • metric_reporters (list) – Список классов для использования в качестве репортеров метрик. Реализация интерфейса AbstractMetricsReporter позволяет подключать классы, которые будут получать уведомления о создании новых метрик. По умолчанию: []
  • metrics_num_samples (int) – Количество образцов для вычисления метрик. По умолчанию: 2
  • metrics_sample_window_ms (int) – Максимальное  количество миллисекунд, используемое для вычисления метрик. По умолчанию: 30000
  • selector (selectors.BaseSelector) – Предоставьте конкретную реализацию селектора для использования при мультиплексировании ввода/вывода. По умолчанию: selectors.DefaultSelector
  • exclude_internal_topics (bool) – определяет, должны ли записи из внутренних топиков(например, смещения) быть доступны консюмеру. Если установлено значение True, то единственным способом получения записей из внутреннего топика будет подписка на него. Требуется 0.10+.  По умолчанию: True
  • sasl_mechanism (str) – Механизм аутентификации, когда security_protocol настроен на SASL_PLAINTEXT или SASL_SSL. Допустимыми значениями являются: PLAIN, GSSAPI, OAUTHBEARER, SCRAM-SHA-256, SCRAM-SHA-512.
  • sasl_plain_username (str) – имя пользователя, необходимое для аутентификации sasl PLAIN и SCRAM. Требуется, если используемый sasl_mechanism - PLAIN или используется один из механизмов SCRAM.
  • sasl_plain_password (str) – пароль для аутентификации sasl PLAIN и SCRAM. Требуется, если используемый sasl_mechanism - PLAIN или используется один из механизмов SCRAM.
  • sasl_kerberos_service_name (str) – Название сервиса, который должен быть включен в механизм GSSAPI sasl. По умолчанию: ‘kafka’
  • sasl_kerberos_domain_name (str) – наименование домена, который должен быть использован в механизме GSSAPI sasl. По умолчанию: один из серверов bootstrap

 

ПРИМЕЧАНИЕ:

Подробное описание параметров настройки доступно в https://kafka.apache.org/documentation/#consumerconfigs

 

assign(partitions)

Присвоение TopicPartitions консюмеров вручную.

Параметры:

partitions (list of TopicPartition) – назначение для данного экземпляра

Указывает:

  • IllegalStateError – если консюмер уже вызвал subscribe().

 

ВНИМАНИЕ!

Невозможно одновременно использовать назначение партиций с помощью  assign() и назначение группы с помощью subscribe().

 

ПРИМЕЧАНИЕ:

Данный интерфейс не поддерживает функцию инкрементального назначения и заменит предыдущее назначение (если таковое имеется).

 

ПРИМЕЧАНИЕ:

Назначение топика вручную с помощью данного метода не затрагивает функционал управлению группой консюмеров. Поэтому при изменении метаданных топика или кластера ребалансировки не будет.

 

assignment()

Получите текущие назначения TopicPartitions для данного консюмера.

Если партиции были назначены напрямую с помощью assign(), данная команда просто вернет те же партиции, что были назначены ранее. Если топики были определены с помощью subscribe(), эта команда выдаст набор партиций топика, назначенные консюмеру (в случае, если назначения еще не было или оно пока еще в процессе, выводом команды будет None).

Возвращает:

{TopicPartition, …}

Тип возврата:

set

 

beginning_offsets(partitions)

Получает первое смещение для заданных разделов.

Этот метод не изменяет текущую потребительскую позицию разделов.

 

ПРИМЕЧАНИЕ:

Если раздел не существует, этот метод может блокироваться до бесконечности.

Параметры:

partitions (list) – Список экземпляров TopicPartition для получения смещений.

Вывод:

int}``: Самые ранние доступные смещения для заданных разделов.

Тип вывода

``{TopicPartition

Указывает:

UnsupportedVersionError - Если брокер не поддерживает поиск смещений по метке времени.

KafkaTimeoutError – Если выборка не удалась в течение request_timeout_ms.

 

bootstrap_connected()

Если бутстрап подключен, возвращает True.

 

close(autocommit=True)

Закрывает консюмера до тех пор, пока не будет произведена необходимая очистка.

Ключевые аргументы:

 

autocommit (bool) – Если для этого консюмера настроен автокоммит, этот флаг заставляет его пытаться зафиксировать смещения до закрытия. По умолчанию: True

 

commit(offsets=None)

Фиксирует смещения в kafka, блокируя их до достижения успеха или получения ошибки.

Это фиксирует смещения только в Kafka. Смещения, зафиксированные с помощью этого API, будут использоваться при первой выборке после каждого ребаланса, а также при запуске. Поэтому, если Вам нужно хранить смещения в чем-то другом (не в Kafka), этот API использовать не следует. Чтобы избежать повторной обработки последнего прочитанного сообщения при перезапуске консюмера, зафиксированное смещение должно быть представлено в виде сообщения, которое должно потреблять Ваше приложение, т. е.: last_offset + 1.

Блокируется до тех пор, пока не завершится фиксация, либо пока не возникнет неустранимая ошибка (в этом случае она будет отброшена вызывающей стороне).

В настоящее время поддерживает только хранение смещений в kafka-топиках (не в zookeeper).

Параметры:

offsets (dict, optional) – {TopicPartition: OffsetAndMetadata}   для  фиксации group_id. По умолчанию используются текущие смещения для всех подписанных разделов.

 

commit_async(offsets=None, callback=None)

Асинхронно фиксирует смещения в kafka, выполняя обратный вызов.

Это фиксирует смещения только в Kafka. Смещения, зафиксированные с помощью этого API, будут использоваться при первой выборке после каждого ребаланса, а также при запуске. Поэтому, если Вам нужно хранить смещения в чем-то другом (не в Kafka), этот API использовать не следует. Чтобы избежать повторной обработки последнего прочитанного сообщения при перезапуске консюмера, зафиксированное смещение должно быть представлено в виде сообщения, которое должно потреблять Ваше приложение, т. е.: last_offset + 1.

Это асинхронный вызов, поэтому блокироваться он не будет. Любые возникшие ошибки либо передаются обратному вызову (если он предусмотрен), либо отбрасываются:

Параметры:

  • offsets (dict, по желанию) – {TopicPartition: OffsetAndMetadata}  используется для фиксации с настроенным group_id. По умолчанию используются текущие смещения для всех подписанных разделов.
  • callback (вызываемый, по желанию) – Вызывается как callback(offsets, response) с ответом в виде исключения или структуры OffsetCommitResponse. Этот обратный вызов может быть использован для запуска пользовательских действий при завершении запроса фиксации.

Выводы:

kafka.future.Future

 

committed(partition, metadata=False)

Получает последнее зафиксированное смещение для данного раздела.

Это смещение будет использоваться в качестве позиции для консюмера в случае сбоя.

Этот вызов может заблокироваться для выполнения удаленного вызова, если данный раздел не назначен этому консюмеру или если консюмер еще не инициализировал  кэширование зафиксированных смещений.

Параметры:

  • partition (TopicPartition) – Раздел, который нужно проверить.
  • metadata (bool, optional) – Если значение = True, то вместо значения offset int возвращается значение OffsetAndMetadata struct. По умолчанию: False.

Выводы:

Последнее зафиксированное смещение (int или OffsetAndMetadata), или None (если предыдущей фиксации не было).

 

end_offsets(partitions)

 

Получает последнее смещение для заданных разделов. Последнее смещение раздела - это смещение предстоящего сообщения, т.е. смещение последнего доступного сообщения + 1.

Этот метод не изменяет текущую позицию разделов.

 

ПРИМЕЧАНИЕ:

Если раздел не существует, этот метод может блокироваться до бесконечности.

Параметры:

partitions (list) – Список экземпляров TopicPartition для получения смещений.

Выводы:

int}``: Конечные смещения для заданных разделов.

Тип выводы:

``{TopicPartition

Указывает:

UnsupportedVersionError - Если брокер не поддерживает поиск смещений по метке времени.

KafkaTimeoutError – Если выборка не удалась в течение request_timeout_ms

 

highwater(partition)

 

Последнее известное highwater смещение для раздела.

Highwater смещение - это смещение, которое будет присвоено следующему создаваемому сообщению. Оно может быть полезно для вычисления отставания. Обратите внимание на то, что и позиция, и highwater относятся к следующему смещению - т. е. highwater смещение на единицу больше, чем самое новое доступное сообщение.

Highwater cмещение возвращается в сообщениях FetchResponse, поэтому оно будет недоступно, если для этого раздела еще не было отправлено ни одного запроса FetchRequests.

Праметры:

partition (TopicPartition) – Раздел, который должен быть проверен

Выводы:

Смещение недоступно

Тип вывода:

int или None

 

metrics(raw=False)

 

Выдает метрики производительности консюмера.

Более подробная информация доступна по ссылке: https://kafka.apache.org/documentation/#consumer_monitoring

 

ВНИМАНИЕ!

Это нестабильный интерфейс. Он может измениться в будущих выпусках без каких-либо предупреждений.

 

offsets_for_times(timestamps)

 

Поиск смещений для заданных разделов по метке времени. Возвращаемое смещение для каждого раздела - это самое раннее смещение, временная метка которого больше или равна заданной временной метке в соответствующем разделе.

Это блокирующий вызов. Консюмер не должен быть назначен разделам.

Если версия формата сообщения в разделе до 0.10.0, т. е. сообщения не имеют временных меток, для этого раздела будет возвращено None. None также будет возвращено для раздела, если в нем нет сообщений..

 

ПРИМЕЧАНИЕ:

Если раздел не существует, этот метод может блокироваться до бесконечности.

 

Парметры:

timestamps (dict) – {TopicPartition: int}  отображение раздела на временную метку. Единицей измерения должны быть миллисекунды с начала определенного события (полночь 1 января 1970 года (UTC)).

Выводы:

OffsetAndTimestamp}``: отображение из раздела на временную метку и смещение первого сообщения с временной меткой, большей или равной целевой временной метке.

Тип вывода:

``{TopicPartition

Указывает:

  • ValueError Если целевая временная метка отрицательна
  • UnsupportedVersionError - Если брокер не поддерживает поиск смещений по метке времени.
  • KafkaTimeoutError – Если в течение request_timeout_ms выборка не удалась

 

partitions_for_topic(topic)

Этот метод сначала проверяет кэш метаданных на наличие информации о топике. Если топик не найден (либо потому, что тема не существует, либо пользователь не авторизован для просмотра топика, либо кэш метаданных не заполнен), то будет выдан вызов обновления метаданных в кластере.

Параметры:

topic (str) – топик, который должен быть проверен.

Выводы:

Id разделов

Тип вывода:

set

 

pause(*partitions)

 

Приостановливает выборку из запрошенных разделов.

Последующие вызовы poll() не будут возвращать никаких записей из этих разделов, пока они не будут возобновлены с помощью resume().

Примечание: Этот метод не влияет на подписку на разделы. В частности, он не вызывает ребалансировку группы, если используется автоматическое назначение.

Параметры:

*partitions (TopicPartition) - раздел, который должен быть проверен.

 

paused()

 

Получите разделы, которые ранее были приостановлены с помощью функции pause().

Выводы:

{partition (TopicPartition), …}

Тип вывода:

set

 

poll(timeout_ms=0, max_records=None, update_offsets=True)

 

Получение данных из назначенных тем/разделов.

Записи извлекаются и возвращаются партиями по топикам и разделам. При каждом опросе консюмер будет пытаться использовать последнее использованное смещение в качестве начального смещения и выполнять выборку последовательно. Последнее потребленное смещение может быть задано вручную через seek() или автоматически как последнее зафиксированное смещение для подписанного списка разделов.

Несовместим с интерфейсом итератора - используйте одно или другое, но не оба сразу.

Параметры:

timeout_ms (int, по желанию) – Миллисекунды, потраченные в ходе опроса в случае, если данные в буфере недоступны. Если 0, возвращается немедленно с любыми записями, которые доступны в данный момент в буфере, иначе возвращается пусто. Не должно быть отрицательным. По умолчанию: 0

max_records (int, по желанию) – Максимальное количество записей, возвращаемых за один вызов poll(). По умолчанию: наследует значение из max_poll_records.

Выводы:

список записей с момента последней выборки для подписного списка топиков и разделов.

Тип вывода:

dict

 

position(partition)

 

Получение смещения следующей записи, которая будет извлечена

Parameters:

partition (TopicPartition) – радел, который должен быть проверен

Returns:

смещение

Return type:

int

 

resume(*partitions)

 

Возобновление выборки из указанных (приостановленных) разделов.

Параметры:

*partitions (TopicPartition) – разделы, которые должны быть резюмированы.

 

seek(partition, offset)

 

Ручное указание смещения выборки для TopicPartition.

Переопределяет смещение выборки, которое консюмер будет использовать при следующем опросе(). Если этот API вызывается для одного и того же раздела несколько раз, при следующем опросе() будет использоваться последнее смещение.

Примечание: Вы можете потерять данные, если этот API будет произвольно использован в середине потребления для сброса смещений выборки.

Параметры:

  • partition (TopicPartition) – Раздел для выполнения операции поиска
  • offset (int) – Смещение сообщения в разделе

Указывают:

AssertionError – Если смещение не является int >= 0; или если раздел в данный момент не назначен.

 

seek_to_beginning(*partitions)

Обращение к самому старому доступному смещению для разделов.

Параметры:

*partitions –укажите конкретные TopicPartitions, в противном случае по умолчанию будут указаны все назначенные разделы.

Указывают:

AssertionError – Если какой-либо раздел в данный момент не назначен, или если не назначен ни один раздел.

 

seek_to_end(*partitions)

 

Обращение к последнему доступному смещению для разделов.

Парметры:

*partitions – укажите конкретные TopicPartitions, в противном случае по умолчанию будут указаны все назначенные разделы

Raises:

AssertionError – Если какой-либо раздел в данный момент не назначен, или если не назначен ни один раздел.

 

subscribe(topics=(), pattern=None, listener=None)

Подпишитесь на список топиков или на regex-шаблонов топика.

Разделы будут динамически назначаться через координатора группы. Подписка на топики не является инкрементной: этот список заменит текущее назначение (если оно есть).

Имейте в виду, что этот метод несовместим с assign().

Параметры:

  • topics (list) – перечень  топиков, на которые нужно подписаться. \
  • pattern (str) - паттерн для сопоставления с имеющимися топиками. Вы должны указать либо топики, либо паттерн, но не оба варианта.
  • listener (ConsumerRebalanceListener) – включите обратный вызов слушателя, который будет осуществляться до и после каждой операции ребалансировки

В рамках управления группами консюмер отслеживает список консюмеров, принадлежащих к определенной группе, и запускает операцию ребалансировки, как толькр происходит одно из следующих событий:

  • Изменение количества разделов для любого из топиков
  • Создается или удаляется топик
  • «Умирает» существующий участник группы консюмеров
  • В группу консюмеров приходит новый участник

При наступлении любого из этих событий предоставленный слушатель будет вызван для указания того, что назначение консюмера было отозвано. При этом есть гарантия того, что разделы, отозванные/назначенные через этот интерфейс, относятся к топикам, подписанным в этом вызове.

Указывает:

  • IllegalStateError- если вызван после предварительного вызова assign().
  • AssertionError – если не указан топик или паттерн.
  • TypeError – если слушатель не является ConsumerRebalanceListener.

 

subscription()

Получить подписку на текущий топик.

Выводы:

{topic, …}

Тип вывода:

set

 

topics()

Получение всех тем, которые консюмер имеет право просматривать. При этом для получения самой последней информации всегда будет выполняться удаленный вызов кластера.

Выводы:

topics

Тип вывода:

set

 

unsubscribe()

Отписка от всех топиков и очистка всех соответствующих разделов.

 

KafkaProducer

classkafka.KafkaProducer(**configs)

Клиент Kafka, который публикует записи в кластере Kafka.

Продюсер безопасен для потоков, использование одного экземпляра продюсера эффективнее, чем использование нескольких экземпляров.

Продюсер работает с пулом буферного пространства, в котором хранятся записи, еще не переданные на сервер, а также с фоновым потоком операций ввода-вывода, который отвечает за превращение этих записей в запросы и передачу их в кластер.

send() является асинхронной. При вызове она добавляет запись в буфер ожидающих отправки записей и немедленно возвращается, что  позволяет консюмеру объединять отдельные записи для повышения эффективности.

Конфиг 'acks' управляет критериями, по которым запросы считаются завершенными. Настройка «all» приводит к блокировке при полной фиксации записи, что является самым медленным, но при этом наиболее долговечным вариантом.

Если запрос не выполняется, продюсер может автоматически повторить попытку (если только 'retries' не настроено на 0). Включение повторных попыток также открывает возможность дублирования (подробнее см. документацию по семантике доставки сообщений: https://kafka.apache.org/documentation.html#semantics ).

Продюсер хранит буферы неотправленных записей для каждого раздела. Эти буферы имеют размер, определяемый конфигурацией 'batch_size'. Увеличение этого размера может привести к большему количеству отправлений, но при этом требуется и гораздо больше памяти (так как обычно у нас есть только один подобный буфер для каждого активного раздела).

По умолчанию буфер доступен для немедленной отправки, даже если в нем есть дополнительное неиспользованное место. Однако если Вы хотите уменьшить количество запросов, то можете установить значение 'linger_ms' больше 0. Это даст указание продюсеру подождать некоторое время и только потом отправить запрос в надежде на то, что придет больше записей, позволяющих заполнить ту же партию. Это похоже на алгоритм Нагла в TCP. Обратите внимание на то, что записи, которые приходят близко друг к другу по времени, обычно собираются в пакет даже при linger_ms=0, поэтому при большой нагрузке пакетная отправка будет происходить независимо от конфигурации linger; однако установка значения больше 0 может привести к уменьшению количества эффективных запросов при отсутствии максимальной нагрузки за счет небольшого увеличения задержки.

Buffer_memory контролирует общий объем памяти, доступный продуюсеру. Если записи отправляются быстрее, чем они могут быть переданы на сервер, то это буферное пространство будет быстро заполнено. Когда буферное пространство будет полностью заполнено, дополнительные вызовы отправки будут автоматически блокироваться.

Key_serializer и value_serializer указывают, как превратить объекты ключа и значения, предоставленные пользователем, в байты.

Ключевые аргумерты:

 

  • bootstrap_servers – ‘host[:port]’ string (or list of ‘host[:port]’ strings) к которому для загрузки начальных метаданных кластера должен обратиться продюсер. Это не обязательно должен быть полный список узлов. В нем должен быть хотя бы один брокер, который ответит на запрос Metadata API. По умолчанию используется порт 9092. Если серверы не указаны, по умолчанию будет использоваться localhost:9092.
  • client_id (str) – имя для этого клиента. Эта строка передается в каждом запросе к серверам и может быть использована для определения конкретных записей журнала на стороне сервера, которые соответствуют этому клиенту. По умолчанию: 'kafka-python-producer-#' (добавляется уникальный номер для каждого экземпляра)
  • key_serializer (вызываемый) – используется для преобразования ключей, предоставленных консюмером, в байты Если не None, вызывается как f(key), должен возвращать байты. По умолчанию: None.
  • value_serializer (callable) – используется для преобразования значений сообщений, предоставленных консюмером, в байты. Если не None, вызывается как f(value) и возвращает байты. По умолчанию: None.
  • acks (0, 1, 'all') –

 

Количество подтверждений, которое прод.сер должен получить от лидера, прежде чем считать запрос завершенным. Это позволяет контролировать долговечность отправляемых записей. Обычно используются следующие настройки:

 

0: прод.сер не будет ждать подтверждений от сервера.

Сообщение будет немедленно добавлено в буфер сокета и будет считаться отправленным. В этом случае нет гарантии, что сервер получил запись, конфигурация повторных попыток не будет действовать (поскольку клиент обычно не знает о неудачах). Смещение, возвращаемое для каждой записи, всегда будет равно -1.

 

1: Дожидается, пока лидер не запишет запись в свой локальный журнал.

Брокер ответит, не дожидаясь полного подтверждения от всех последователей. В этом случае, если лидер выйдет из строя сразу после подтверждения записи, но до того, как ее воспроизведут последователи, запись будет потеряна.

all: Дожидается того момента, когда полный набор синхронизированных реплик запишет запись.

Это гарантирует, что запись не будет потеряна до тех пор, пока «жива» хотя бы одна синхронизированная реплика. Это самая надежная из всех доступных гарантий. Если значение не установлено, по умолчанию используется acks=1.

  • compression_type (str) – Тип сжатия для всех данных, генерируемых продюсером. Допустимые значения: 'gzip', 'snappy', 'lz4' или None. Сжатие происходит из полных пакетов данных, поэтому эффективность пакетной обработки влияет и на степень сжатия (большая степень пакетной обработки означает лучшее сжатие). По умолчанию: None.

 

  • retries (int) – Если установить значение больше нуля, клиент будет повторно отправлять любую запись, отправка которой не удалась ввиду возникшей ошибки. Обратите внимание на то, что эта повторная попытка ничем не отличается от того, если бы клиент повторно отправил запись после получения ошибки. Разрешение повторных попыток без установки max_in_flight_requests_per_connection в 1 потенциально изменит порядок записей, поскольку если две партии отправляются в один раздел, и первая не удалась и была повторно отправлена, а вторая удалась, то записи из второй партии могут появиться первыми. По умолчанию: 0.

 

 

  • batch_size (int) – Запросы, отправляемые брокерам, будут содержать несколько пакетов, по одному на каждый раздел с данными. Небольшой размер партии может снизить пропускную способность отправки этих запросов (размер партии, равный нулю, полностью отключит отправку). По умолчанию: 16384

 

  • linger_ms (int) – Продюсер объединяет все записи, поступившие в промежутке между передачами запросов, в один пакетный запрос. Обычно это происходит только тогда, когда записи поступают быстрее, чем их можно отправить. Однако при некоторых обстоятельствах клиент может захотеть уменьшить количество запросов даже при умеренной нагрузке. Это достигается путем добавления небольшой искусственной задержки; то есть вместо немедленной отправки записи продюсер будет ждать заданной задержки для того, чтобы позволить другим записям быть отправленными. Это можно представить как аналог алгоритма Нагла в TCP. Эта настройка задает верхнюю границу задержки для пакетной отправки: как только мы получим пакетное количество записей для раздела, оно будет отправлено немедленно, независимо от этой настройки, однако если у нас накопилось меньше, чем это количество байт для этого раздела, мы «задержимся» на указанное время, ожидая появления новых записей. По умолчанию эта настройка равна 0 (т. е. задержки нет). Установка linger_ms=5 приведет к уменьшению количества отправляемых запросов, но добавит до 5 мс задержки к отправляемым записям при отсутствии нагрузки. По умолчанию: 0.

 

  • partitioner (вызываемый) – Вызываемый модуль, используемый для определения того, к какому разделу относится каждое сообщение. Вызывается следующим образом (после сериализации ключа): partitioner(key_bytes, all_partitions, available_partitions). Реализация partitioner по умолчанию хэширует каждый не-None ключ, используя тот же алгоритм murmur2, что и java-клиент, так что сообщения с одинаковым ключом назначаются одному и тому же разделу. Если ключ = None, сообщение доставляется в случайный раздел (по возможности отфильтрованный только в разделы с доступными лидерами).

 

  • buffer_memory (int) – Общее количество байт памяти, которое продюсер должен использовать для буферизации записей, ожидающих отправки на сервер. Если записи отправляются быстрее, чем они могут быть доставлены на сервер, продюсер будет блокировать их до max_block_ms. В текущей реализации это значение является приблизительным. По умолчанию: 33554432 (32 МБ)

 

  • connections_max_idle_ms – Закрывает простаивающие соединения через количество миллисекунд, указанное в этом конфиге. Брокер закрывает простаивающие соединения после connections.max.idle.ms, что позволяет избежать неожиданных ошибок отключения сокетов на клиенте. По умолчанию: 540000

 

  • max_block_ms (int) – Количество миллисекунд для блокировки во время send() и partitions_for(). Эти методы могут быть заблокированы либо из-за переполнения буфера, либо из-за недоступности метаданных. Блокировка в пользовательских сериализаторах или разделителях не будет учитываться в этом таймауте. По умолчанию: 60000.

 

  • max_request_size (int) – Максимальный размер запроса. Это также фактически ограничение на максимальный размер записи. Обратите внимание, что сервер имеет свое собственное ограничение на размер записи, которое может отличаться от этого. Эта настройка ограничивает количество пакетов записей, которые производитель будет отправлять в одном запросе, чтобы избежать отправки огромных запросов. По умолчанию: 1048576.

 

  • metadata_max_age_ms (int) – Период времени в миллисекундах, по истечении которого мы принудительно обновляем метаданные для того, чтобы проактивно обнаружить новые брокеры или разделы. По умолчанию: 300000

 

  • retry_backoff_ms (int) Миллисекунды для повторных попытках при возникновении ошибок. По умолчанию: 100.

 

  • request_timeout_ms (int) – Таймаут запроса клиента в миллисекундах. По умолчанию: 30000.

 

  • receive_buffer_bytes (int) – Размер буфера приема TCP (SO_RCVBUF), который будет использоваться при чтении данных. По умолчанию: None (полагается на системные настройки по умолчанию). Java-клиент использует по умолчанию 32768

 

  • send_buffer_bytes (int) Размер буфера отправки TCP (SO_SNDBUF), который будет использоваться при отправке данных. По умолчанию: None (полагается на системные настройки по умолчанию). Java-клиент по умолчанию принимает значение 131072

 

  • socket_options (list) – Список кортежей-аргументов для socket.setsockopt, применяемых к сокетам брокерских соединений. По умолчанию: [(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)].

 

  • reconnect_backoff_ms (int) – Количество времени в миллисекундах, которое необходимо выждать перед попыткой повторного подключения к данному узлу. По умолчанию: 50.

 

  • reconnect_backoff_max_ms (int) – Максимальное количество времени в миллисекундах для отката/ожидания при повторном подключении к брокеру, который неоднократно не смог установить соединение. Если указано, то время ожидания для каждого узла будет экспоненциально увеличиваться для каждого последовательного сбоя соединения, вплоть до этого максимума. После достижения максимума попытки переподключения будут продолжаться периодически с этой фиксированной скоростью. Во избежание шторма соединений к коэффициенту обратного хода будет применен коэффициент рандомизации 0,2, в результате чего он будет находиться в случайном диапазоне между 20 % ниже и 20 % выше вычисленного значения. По умолчанию: 1000.

 

  • max_in_flight_requests_per_connection (int) - Запросы передаются брокерам kafka до этого максимального количества запросов на одно брокерское соединение. Обратите внимание на  то, что если этот параметр установлен больше 1 и есть неудачные отправки, существует риск переупорядочивания сообщений из-за повторных попыток (т. е. если повторные попытки включены). По умолчанию: 5.

 

  • security_protocol (str) – Протокол, используемый для связи с брокерами. Допустимыми значениями являются: PLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL. По умолчанию: PLAINTEXT.

 

  • ssl_context (ssl.SSLContext) – предварительно сконфигурированный SSLContext для обертывания сокетных соединений. Если указано, все остальные конфигурации ssl_* будут проигнорированы. По умолчанию: None.

 

  • ssl_check_hostname (bool) – флаг для настройки того, должен ли ssl_handshake проверять соответствие сертификата имени хоста брокера. По умолчанию: true.

 

  • ssl_cafile (str) – необязательное имя файла ca для использования при проверке сертификата. По умолчанию: none.

 

  • ssl_certfile (str) – необязательное имя файла в формате pem, содержащего сертификат клиента, а также любые сертификаты ca, необходимые для установления подлинности сертификата. По умолчанию: none.

 

  • ssl_keyfile (str) – необязательное имя файла, содержащего закрытый ключ клиента. По умолчанию: none.

 

  • ssl_password (str) – необязательный пароль, который будет использоваться при загрузке цепочки сертификатов. По умолчанию: none.

 

  • ssl_crlfile (str) – необязательное имя файла, содержащего CRL для проверки срока действия сертификата. По умолчанию проверка CRL не производится. При указании файла только лист сертификата будет проверен по этому CRL. CRL может быть проверен только с помощью Python 3.4+ или 2.7.9+. По умолчанию: none.

 

  • ssl_ciphers (str) –задает доступные шифры для ssl-соединений. Это должна быть строка в формате списка шифров OpenSSL. Если ни один шифр не может быть выбран (потому что опции времени компиляции или другая конфигурация запрещают использование всех указанных шифров), будет выдана ошибка ssl.SSLError. См. ssl.SSLContext.set_ciphers

 

  • api_version (tuple) – Укажите, какую версию API Kafka Вы хотите использовать. Если установлено значение None, клиент будет пытаться определить версию брокера путем опроса различных API. Пример: (0, 10, 2). По умолчанию: None.

 

  • api_version_auto_timeout_ms (int) – количество миллисекунд для выброса исключения таймаута из конструктора при проверке версии api брокера. Применяется, только если api_version установлено как None.

 

  • metric_reporters (list) – Список классов для использования в качестве репортеров метрик. Реализация интерфейса AbstractMetricsReporter позволяет подключать классы, которые будут получать уведомления о создании новых метрик. По умолчанию: []

 

  • metrics_num_samples (int) – Количество образцов, используемых для расчета метрик. По умолчанию: 2

 

  • metrics_sample_window_ms (int) – максимальное количество времени в миллисекундах, необходимое для подсчета метрик. По умолчанию: 30000

 

  • selector (selectors.BaseSelector) - Предоставляет конкретную реализацию селектора для использования при мультиплексировании операций ввода/вывода. По умолчанию: Selector.

 

  • sasl_mechanism (str) – Механизм аутентификации, когда security_protocol настроен на SASL_PLAINTEXT или SASL_SSL. Допустимыми значениями являются: PLAIN, GSSAPI, OAUTHBEARER, SCRAM-SHA-256, SCRAM-SHA-512.

 

  • sasl_plain_username (str) - имя пользователя для аутентификации sasl PLAIN и SCRAM. Требуется, если sasl_mechanism - PLAIN или один из механизмов SCRAM.

 

  • sasl_plain_password (str) - пароль для аутентификации sasl PLAIN и SCRAM. Требуется, если sasl_mechanism - PLAIN или один из механизмов SCRAM.

 

  • sasl_kerberos_service_name (str) – Имя службы для включения в рукопожатие механизма GSSAPI sasl. По умолчанию: ‘kafka’

 

  • sasl_kerberos_domain_name (str) – имя домена kerberos для использования в рукопожатии механизма GSSAPI sasl. По умолчанию: один из загрузочных серверов

 

  • sasl_oauth_token_provider (AbstractTokenProvider) – Экземпляр поставщика токенов OAuthBearer. (См. kafka.oauth.abstract). По умолчанию: None

 

ПРИМЕЧАНИЕ:

Более подробная информация качательно параметров настройки доступна по ссылке: https://kafka.apache.org/0100/configuration.html#producerconfigs

 

bootstrap_connected()

Возвращает True, если бутстрап подключен.

 

close(timeout=None)

 

Закрывает этого продюсера.

Параметры:

timeout (float, optional) - таймаут в секундах для ожидания завершения.

 

flush(timeout=None)

 

Вызов этого метода делает все буферизованные записи доступными для отправки (даже если linger_ms больше 0) и блокирует завершение запросов, связанных с этими записями. Пост-условием flush() является то, что любая ранее отправленная запись будет завершена (например, Future.is_done() == True). Запрос считается завершенным, если он либо успешно подтвержден в соответствии с конфигурацией 'acks' для производителя, либо привел к ошибке.

Другие потоки могут продолжать отправлять сообщения, пока один поток заблокирован в ожидании завершения вызова flush; однако завершение сообщений, отправленных после начала вызова flush, не гарантируется..

Параметры:

timeout (float, optional) – таймаут в секундах для ожидания завершения.

Указывает:

KafkaTimeoutError – неспособность просмотреть буферизованные записи в течение заданного таймаута

 

metrics(raw=False)

 

Получает метрики производительности работы продюсера.

Взято из Java- Producer, более подробная информация доступна по ссылке: https://kafka.apache.org/documentation/#producer_monitoring

 

ВНИМАНИЕ

Это нестабильный интерфейс. Он может измениться в более поздних выпусках без специального предупреждения.

 

partitions_for(topic)

 

Возвращает набор всех известных разделов для данного топика.

 

send(topic, value=None, key=None, headers=None, partition=None, timestamp_ms=None)

 

Публикует сообщение в топик.

Параметры

  • topic (str) – топик, в котором будет опубликовано сообщение
  • value (optional) – значение сообщения. Должно иметь тип bytes или быть сериализуемым в bytes через настроенный value_serializer. Если значение равно None, требуется ключ, а сообщение действует как 'delete'. Смотрите официальную инструкцию kafka по сжатию: https://kafka.apache.org/documentation.html#compaction (для сжатия нужна версия kafka >= 0.8.1)
  • partition (int, optional) –указыввет раздел. Если значение не задано, раздел будет выбран с помощью предварительно настроенного «разделителя».
  • key (optional) – ключ, подобранный к сообщению. Может использоваться для определения того,  в какой раздел отправить то или иное сообщение. Если раздел = None (и настройки разделителя продсера оставлены такими, какими они были по умолчанию), тогда сообщение с тем же ключом будет доставлен в тот же самый раздел (если ключ= None, разделение выбирается случайным образом).
  • headers (optional) – перечень пар ключевых значений заголовок. Перечисляет наименования кортежей.
  • timestamp_ms (int, optional) – миллисекунды (от 01.01. 1970 UTC) для использования временной метки сообщения.

Выводы:

Относится к RecordMetadata

Тип вывода:

FutureRecordMetadata

Указывает:

KafkaTimeoutError – если невозможно получить метаданные топика или буфер памяти установлен на max_block_ms

 

KafkaAdminClient

classkafka.KafkaAdminClient(**configs)

Класс, предназначенный для управления кластером Kafka.

ВНИМАНИЕ!

Это нестабильный интерфейс, который был добавлен совсем недавно и может быть изменен

изменениям без предупреждения.

В частности, многие методы в настоящее время возвращают необработанные кортежи протоколов. В будущих выпусках мы планируем превратить их в более приятные для работы объекты. К сожалению, это усовершенствование, скорее всего, приведет к поломке интерфейсов.

Класс KafkaAdminClient будет вести переговоры о получении последней версии каждого формата протокола сообщений, поддерживаемого как клиентской библиотекой kafka-python, так  и брокером Kafka. Использование необязательных полей из версий протоколов, которые не поддерживаются брокером, приведет к появлению исключений класса IncompatibleBrokerVersion.

 

Использование этого класса требует того, что версия брокера была как минимум  0.10.0.0.

 

Ключевые аргументы:

 

  • bootstrap_servers – ‘host[:port]’ Строка (или список строк 'host[:port]'), к которой потребитель должен обратиться для загрузки начальных метаданных кластера. Это не обязательно должен быть полный список узлов. Он просто должен содержать хотя бы одного брокера, который ответит на запрос Metadata API. По умолчанию используется порт 9092. Если серверы не указаны, по умолчанию будет использоваться localhost:9092.
  • client_id (str) - имя для этого клиента. Эта строка передается в каждом запросе к серверам и может быть использована для идентификации конкретных записей журнала на стороне сервера, соответствующих данному клиенту. Также передается GroupCoordinator для ведения журнала в отношении администрирования групп консюмеров. По умолчанию: ‘kafka-python-{version}’
  • reconnect_backoff_ms (int) – Количество времени в миллисекундах до попытки повторного подключения к данному хосту. По умолчанию: 50.
  • reconnect_backoff_max_ms (int) – Максимальное количество времени в миллисекундах для отката/ожидания при повторном подключении к брокеру, который не смог установить соединение. Время ожидания для каждого узла будет увеличиваться экспоненциально при каждом последовательном сбое соединения, вплоть до достижения максимума, после чего попытки переподключения будут продолжаться через определенные промежутки времени. Чтобы избежать «шторма соединений», к коэффициенту обратного хода будет применен коэффициент рандомизации 0,2, что приведет к получению случайного диапазона - 20 % ниже и 20 % выше вычисленного значения. По умолчанию: 1000.
  • request_timeout_ms (int) – Таймаут запроса клиента в миллисекундах. По умолчанию: 30000.
  • connections_max_idle_ms – Закрывает простаивающие соединения через количество миллисекунд, указанное в этом конфиге. Брокер закрывает простаивающие соединения после connections.max.idle.ms, что позволяет избежать неожиданных ошибок отключения сокетов на клиенте. По умолчанию: 540000
  • retry_backoff_ms (int) – Миллисекунды для отката при повторных попытках при ошибках. По умолчанию: 100.
  • max_in_flight_requests_per_connection (int) – Запросы передаются брокерам kafka вплоть до достижения указанного количества максимальных запросов на одно соединение с брокером. По умолчанию: 5.
  • receive_buffer_bytes (int) – Размер буфера приема TCP (SO_RCVBUF), который будет использоваться при чтении данных. По умолчанию: None (полагается на системные настройки по умолчанию). Java-клиент использует по умолчанию 32768.
  • send_buffer_bytes (int) – Размер буфера отправки TCP (SO_SNDBUF), который будет использоваться при отправке данных. По умолчанию: None (полагается на системные настройки по умолчанию). Java-клиент по умолчанию принимает значение 131072.
  • socket_options (list) – Список кортежей-аргументов для socket.setsockopt, применяемых к сокетам брокерских соединений. По умолчанию: [(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)].
  • metadata_max_age_ms (int) – Период времени в миллисекундах, по истечении которого мы принудительно обновляем метаданные (даже если мы не видели никаких изменений в руководстве разделами) для того, чтобы проактивно обнаружить новые брокеры или разделы. По умолчанию: 300000
  • security_protocol (str) – Протокол, используемый для связи с брокерами. Допустимыми значениями являются: PLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL. По умолчанию: PLAINTEXT.
  • ssl_context (ssl.SSLContext) – Предварительно сконфигурированный SSLContext для обертывания сокетных соединений. Если указано точное значение, все остальные конфигурации ssl_* будут проигнорированы. По умолчанию: None.
  • ssl_check_hostname (bool) – Флаг для настройки того, должен ли SSL-кейк проверять соответствие сертификата имени хоста брокера или нет. По умолчанию: True.
  • ssl_cafile (str) – Необязательное имя файла ЦС для использования при проверке сертификата. По умолчанию: None.
  • ssl_certfile (str) – Необязательное имя файла в формате PEM, содержащего сертификат клиента, а также любые сертификаты ЦС, необходимые для установления подлинности сертификата. По умолчанию: None.
  • ssl_keyfile (str) – Необязательное имя файла, содержащего закрытый ключ клиента. По умолчанию: None.
  • ssl_password (str) – Необязательный пароль, который будет использоваться при загрузке цепочки сертификатов. По умолчанию: None.
  • ssl_crlfile (str) – Необязательное имя файла, содержащего CRL для проверки срока действия сертификата. По умолчанию проверка CRL не производится. При указании файла только лист сертификата будет проверен по этому CRL. CRL может быть проверен только с Python 3.4+ или 2.7.9+. По умолчанию: None.
  • api_version (tuple) – Указывает, какую версию API Kafka стоит использовать. Если установлено значение None, KafkaClient будет пытаться определить версию брокера путем опроса различных API. Пример: (0, 10, 2). По умолчанию: None
  • api_version_auto_timeout_ms (int) – количество миллисекунд для выброса исключения таймаута из конструктора при проверке версии api брокера. Применяется только в том случае, если api_version равно None.
  • selector (selectors.BaseSelector) – Предоставьте конкретную реализацию селектора для использования при мультиплексировании ввода/вывода. По умолчанию: selectors.DefaultSelector.
  • metrics (kafka.metrics.Metrics) – указывает экземпляр метрики для сбора статистики сетевого ввода-вывода. По умолчанию: None.
  • metric_group_prefix (str) – префикс для названий метрик. По умолчанию: ‘’
  • sasl_mechanism (str) – Механизм аутентификации, когда security_protocol настроен на SASL_PLAINTEXT или SASL_SSL. Допустимыми значениями являются: PLAIN, GSSAPI, OAUTHBEARER, SCRAM-SHA-256, SCRAM-SHA-512.
  • sasl_plain_username (str) – имя пользователя для аутентификации sasl PLAIN и SCRAM. Требуется, если sasl_mechanism - PLAIN или один из механизмов SCRAM.
  • sasl_plain_password (str) – пароль для аутентификации sasl PLAIN и SCRAM. Требуется, если sasl_mechanism - PLAIN или один из механизмов SCRAM.
  • sasl_kerberos_service_name (str) – Имя сервиса для включения в механизм GSSAPI sasl. По умолчанию: 'kafka'
  • sasl_kerberos_domain_name (str) – имя домена kerberos для использования в рукопожатии механизма GSSAPI sasl. По умолчанию: один из загрузочных серверов
  • sasl_oauth_token_provider (AbstractTokenProvider) – Экземпляр поставщика токенов OAuthBearer. (См. kafka.oauth.abstract). По умолчанию: None

 

alter_configs(config_resources)

Изменение параметров конфигурации одного или нескольких ресурсов Kafka.

 

ВНИМАНИЕ!

В настоящее время это не работает для ресурсов BROKER, потому что они должны быть отправлены конкретному брокеру, в то время как уэтот параметр всегда выбирает наименее загруженный узел. Подробности см. в комментарии к исходном коду.

Параметры:

config_resources – список всех объектов ConfigResource objects.

Выводы:

Нужная версия класса AlterConfigsResponse.

 

close()

Закрывает соединение KafkaAdminClient с брокером Kafka.

 

create_acls(acls)

Создает список  ACL

Эта конечная точка принимает только список конкретных объектов ACL, без ACLFilters. Выводит ошибку TopicAlreadyExistsError в том случае, если такой топик уже присутствует.

Параметры:

acls – список объектов ACL

Выводы:

Сообщает об успешном выполнении операции/неудачной попытке

 

create_partitions(topic_partitions, timeout_ms=None, validate_only=False)

Создание дополнительных разделов для существующего топика.

Параметры:

  • topic_partitions – Карта строк названий топиков для объектов NewPartition.
  • timeout_ms – Миллисекунды для ожидания создания новых разделов перед возвратом брокера.
  • validate_only – Если значение True, то создавать новые разделы не нужно. По умолчанию: False

Выводы:

Необходимая версия класса CreatePartitionsResponse.

 

create_topics(new_topics, timeout_ms=None, validate_only=False)

Создает новый топик в кластере.

Параметры:

  • new_topics – список объектов NewTopic.
  • timeout_ms – Миллисекунды ожидания создания новых топиков перед их возвратом брокера.
  • validate_only – Если значение True, ни в коем случае не создавайте новые топики. Поддерживается не всеми версиями. По умолчанию: False

Выводы:

Необходимая версия класса CreateTopicResponse

 

delete_acls(acl_filters)

Удаляет набор ACL

Удаляет все ACL, соответствующие списку входных ACLFilter

Параметры:

  • acl_filters – список ACLFilter

Выводы:

Список, состоящий из 3 кортежей, соответствующих списку входных фильтров. Кортежи содержат (входной ACLFilter, список затронутых ACL, а также экземпляр KafkaError)

 

delete_topics(topics, timeout_ms=None)

Удаляет топики из кластера.

Параметры:

  • topics – Список строк, содержащих название топиков.
  • timeout_ms – Миллисекунды до удаления топиков перед возвращением брокера.

Выводы:

Необходимая версия класса DeleteTopicsResponse.

 

describe_acls(acl_filter)

Описывает набо ACL

Используется для возврата набора ACL, соответствующих заданному фильтру ACLFilter. Для этого кластер должен быть сконфигурирован с авторизатором, иначе Вы получите ошибку SecurityDisabledError

Параметры:

  • acl_filter – объект ACLFilter

Возвраты:

кортеж из списка подходящих объектов ACL и ошибки KafkaError (NoError в случае успеха)

 

describe_configs(config_resources, include_synonyms=False)

Получение параметров конфигурации для одного или нескольких ресурсов Kafka.

Параметры:

  • config_resources – Список объектов ConfigResource. Любые ключи в дикте ConfigResource.configs будут использованы для фильтрации результатов. При установке диктанта configs в None будут получены все значения. Пустой dict получит нулевые значения (согласно протоколу Kafka).
  • include_synonyms – Если значение True, то в ответе будут возвращены синонимы. Поддерживается не всеми версиями. По умолчанию: False.

Выводы:

Необходимая версия DescribeConfigsResponse.

 

describe_consumer_groups(group_ids, group_coordinator_id=None, include_authorized_operations=False)

Описывает набор групп консюмеров.

Сообщения о любых ошибках приходят немедленно.

Параметры:

  • group_ids – Список идентификаторов групп консюмеров. Обычно это имена групп в виде строк.
  • group_coordinator_id – Node_id брокера-координатора групп. Если установлено значение None, то для каждой группы будет запрошен кластер для того, чтобы найти координатора этой группы. Явное указание этого параметра позволит  избежать лишних сетевых обходов, если Вы уже знаете координатора группы. Это полезно только в том случае, если все group_ids имеют одного и того же координатора, иначе появится ошибка. По умолчанию: None.
  • include_authorized_operations – определяет,  включать или не включать информацию о том, какие операции разрешено выполнять группе. Поддерживается только в API версии >= v3. По умолчанию: False.

Выводы:

Список описаний групп. На данный момент описания групп - это необработанные результаты из ответа DescribeGroupsResponse. В перспективе мы планируем изменить его для того, чтобы он возвращал именованные кортежи, а также декодировал назначения разделов.

 

list_consumer_group_offsets(group_id, group_coordinator_id=None, partitions=None)

Получение смещений консюмеров для одной группы консюмеров.

Примечание: При этом не проверяется, действительно ли group_id существуют в кластере.

Как только возникает какая-либо ошибка, сразу же приходит соответствующее сообщение.

Параметры:

  • group_id – Имя идентификатора группы потребителей, для которой необходимо получить смещения.
  • group_coordinator_id – Node_id брокера-координатора группы. Если задано значение None, будет сделан запрос в кластер для поиска координатора группы. Явное указание этого параметра может быть полезно для предотвращения лишних сетевых обходов, если вы уже знаете координатора группы. По умолчанию: None.
  • partitions – Список TopicPartitions, для которых необходимо получить смещения. В брокерах >= 0.10.2 это значение может быть установлено в None, чтобы получить все известные смещения для группы потребителей. По умолчанию: None.

Выводы:

Словарь с ключами TopicPartition и значениями OffsetAndMetada. Разделы, которые не указаны и для которых group_id не имеет записанного смещения, опускаются. Значение смещения -1 означает, что у group_id нет смещения для данного TopicPartition. Значение -1 может иметь место только для разделов, которые явно указаны.

 

list_consumer_groups(broker_ids=None)

Список всех групп консюмеров, известных кластеру.

Возвращает список кортежей групп консюмеров. Кортежи состоят из имени группы консюмеров и типа протокола группы консюмеров.

Возвращаются только те группы консюмеров, которые хранят свои смещения в Kafka. Тип протокола будет пустой строкой для групп, созданных с помощью API Kafka < 0.9, поскольку, хотя они и хранят свои смещения в Kafka, они не используют Kafka для координации групп. Для групп, созданных с помощью Kafka >= 0.9, тип протокола обычно будет «consumer».

Как только возникает какая-либо ошибка, сразу же приходит соответствующее сообщение.

Параметры:

  • broker_ids – Список идентификаторов узлов брокера для запроса групп консюмеров. Если установлено значение None, будут запрошены все брокеры в кластере. Явное указание брокера(ов) может быть полезно для определения того, какие группы потребителей координируются этим брокером(ами). По умолчанию: None

Выводы:

Список кортежей групп консюмеров.

Указывает:

  • GroupCoordinatorNotAvailableError – Координатор недоступен, поэтому не может обрабатывать запросы.
  • GroupLoadInProgressError – Координатор загружается и поэтому не может обрабатывать запросы.

 

KafkaClient

classkafka.KafkaClient(**configs)

Сетевой клиент для асинхронных запросов/ответов при вводе/выводе данных по сети.

Это внутренний класс, используемый для реализации пользовательских клиентов продюсера и консюмера.

Этот класс не является потокобезопасным!

 

cluster

Локальный кэш метаданных кластера, получаемый через MetadataRequests во время poll().

Тип:

ClusterMetadata

 

 Ключевые аргументы:

  • bootstrap_servers – ‘host[:port]’ строка(or list of ‘host[:port]’ strings) к которой должен обратиться клиент для загрузки начальных метаданных кластера. Это не обязательно должен быть полный список узлов. В нем должен быть хотя бы один брокер, который ответит на запрос Metadata API. По умолчанию используется порт 9092. Если серверы не указаны, по умолчанию будет использоваться localhost:9092.
  • client_id (str) – имя для этого клиента. Эта строка передается в каждом запросе к серверам и может использоваться для идентификации определенных записей журнала на стороне сервера, которые соответствуют этому клиенту. Также представлен координатору группы для регистрации в отношении администрирования групп потребителей. По умолчанию: «kafka-python- {version}»
  • reconnect_backoff_ms (int) – Время ожидания повторного подключения к данному узлу в миллисекундах. По умолчанию:50.
  • reconnect_backoff_max_ms (int) – Максимальное время (в миллисекундах) задержки/ожидания при повторном подключении к брокеру, которому неоднократно не удавалось подключиться. Если этот параметр установлен, то при каждом последовательном сбое соединения величина задержки на узел будет увеличиваться экспоненциально до этого максимума. После достижения максимального значения попытки повторного подключения будут периодически продолжаться с этой фиксированной скоростью. Чтобы избежать штормов соединения, к задержке будет применяться коэффициент рандомизации 0,2, что приведет к случайному диапазону от 20% ниже до 20% выше расчетного значения. По умолчанию: 1000.
  • request_timeout_ms (int) – Время ожидания запроса клиента в миллисекундах. По умолчанию: 30000.
  • connections_max_idle_ms - Закрытие неиспользуемых подключений по истечении количества миллисекунд, указанного в этом файле конфигурации. Брокер закрывает свободные соединения после connections.max.idle.ms, так что это позволяет избежать попадания неожиданных отключенных ошибок сокета на клиенте. По умолчанию: 540000
  • retry_backoff_ms (int) – Миллисекунды для отката при повторении ошибок. По умолчанию: 100.
  • max_in_flight_requests_per_connection (int) – Запросы передаются брокерам kafka до этого количества максимальных запросов на соединение с брокером. По умолчанию: 5.
  • receive_buffer_bytes (int) – Размер буфера приема TCP (SO_RCVBUF), используемого при чтении данных. По умолчанию: None (зависит от системных настроек по умолчанию). Клиент Java по умолчанию имеет значение 32768.
  • send_buffer_bytes (int) – Размер буфера передачи TCP (SO_SNDBUF), используемого при передаче данных. По умолчанию: None (зависит от системных настроек по умолчанию). Клиент Java по умолчанию имеет значение 131072.
  • socket_options (list) - Список кортежей-аргументов к socket.setsockopt для применения к сокетам подключения брокера. По умолчанию: [(socket.IPPROTO_TCP, сокет. TCP_NODELAY, 1)]
  • metadata_max_age_ms (int) – Период времени в миллисекундах, после которого мы принудительно обновляем метаданные, даже если мы не видели никаких изменений в лидерстве разделов, чтобы проактивно обнаруживать какие-либо новые брокеры или разделы. По умолчанию: 300000
  • security_protocol (str) – Протокол, используемый для связи с брокерами. Допустимые значения: PLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL. По умолчанию: PLAINTEXT.
  • ssl_context (ssl.SSLContext) – Предварительно настроенный SSLContext для упаковки соединений сокетов. Если указано, все остальные конфигурации ssl _ * будут игнорироваться. По умолчанию: None.
  • ssl_check_hostname (bool) – Флаг, указывающий, должно ли SSL-квитирование проверять соответствие сертификата имени узла брокера. По умолчанию: True.
  • ssl_cafile (str) – Необязательное имя файла ЦС для использования при проверке сертификата. По умолчанию: None.
  • ssl_certfile (str) – Необязательное имя файла в формате PEM, содержащего сертификат клиента, а также любые сертификаты ЦС, необходимые для установления подлинности сертификата. По умолчанию: None.
  • ssl_keyfile (str) – Необязательное имя файла, содержащее закрытый ключ клиента. По умолчанию: None.
  • ssl_password (str) – Необязательный пароль для загрузки цепочки сертификатов. По умолчанию: None.
  • ssl_crlfile (str) – Необязательное имя файла, содержащее CRL для проверки срока действия сертификата. По умолчанию проверка CRL не выполняется. При предоставлении файла по этому CRL будет проверяться только конечный сертификат. CRL можно проверить только с помощью Python 3.4 + или 2.7.9 +. По умолчанию: None.
  • ssl_ciphers (str) – устанавливает доступные шифры для ssl-соединений. Это должна быть строка в формате списка шифров OpenSSL. Если ни один шифр не может быть выбран (поскольку параметры compile-time или другая конфигурация запрещают использование всех указанных шифров), ssl. Будет поднят SSLError. Смотрите ssl. SSLContext.set_ciphers
  • api_version (tuple) – Указывает, какую версию API Kafka нужно использовать. Если установлено значение None, KafkaClient попытается вывести версию брокера путем зондирования различных API. Пример: (0, 10, 2). По умолчанию: None
  • api_version_auto_timeout_ms (int) – количество миллисекунд, необходимое для выдачи исключения тайм-аута из конструктора при проверке версии API брокера. Применяется только в том случае, если api_version имеет значение None
  • selector (selectors.BaseSelector) – Предоставьте конкретную реализацию селектора для мультиплексирования ввода-вывода. По умолчанию: селекторы. Селектор по умолчанию.
  • metrics (kafka.metrics.Metrics) – Дополнительно предоставляет экземпляр метрики для сбора статистики ввода-вывода сети. По умолчанию: None.
  • metric_group_prefix (str) – Префикс для имен метрик. По умолчанию: ‘’
  • sasl_mechanism (str) – Механизм аутентификации, когда security_protocol настроен для SASL_PLAINTEXT или SASL_SSL. Допустимые значения: PLAIN, GSSAPI, OAUTHBEARER, SCRAM-SHA-256, SCRAM-SHA-512.
  • sasl_plain_username (str) – имя пользователя для аутентификации sasl PLAIN и SCRAM. Требуется, если sasl_mechanism является PLAIN или одним из механизмов SCRAM.
  • sasl_plain_password (str) – пароль для аутентификации sasl PLAIN и SCRAM. Требуется в случае, если sasl_mechanism является PLAIN или одним из механизмов SCRAM.
  • sasl_kerberos_service_name (str) – Имя службы, включаемое в квитирование механизма sasl GSSAPI. По умолчанию: «kafka»
  • sasl_kerberos_domain_name (str) – доменное имя kerberos для использования при квитировании механизма GSSAPI sasl. По умолчанию: один из серверов начальной загрузки
  • sasl_oauth_token_provider (AbstractTokenProvider) экземпляр поставщика маркеров OAuthBearer. (См. Kafka.oauth.abstract). По умолчанию: None

 

add_topic(topic)

 

Добавление топика в список топиков, отслеживаемых с помощью метаданных.

Параметры:

topic (str) – топик, который необходимо отслеживать

Выводы:

появляется после запроса/ответа метаданных

Тип выводы:

Future

 

bootstrap_connected()

 

Если подключен узел начальной загрузки, возвращает True

 

check_version(node_id=None, timeout=2, strict=False)

 

Попытка угадать версию брокера Kafka.

Примечание: Возможно, что использование этого метода занимает большего времени, чем хотелось бы ( превышает указанное время ожидания). Это может произойти, если весь кластер не работает и клиент переходит в режим задержки загрузки. Это возможно только в том случае, если node_id имеет значение None.

Возвращает: кортеж версии, т.е. (0, 10), (0, 9), (0, 8, 2),...

Указывает:

  • NodeNotReadyError (при наличии node_id)
  • NoBrokersAvailable (если node_id отсутствует)
  • UncognizedBrokerVersion - пожалуйста, зафиксируйте ошибку, если она была обнаружена!
  • AssertionError (если strict = True) - пожалуйста, зафиксируйте ошибку, если она была обнаружена!

 

close(node_id=None)

 

Закрывает одно или все брокерские соединения.

Параметры:

node_id (int, optional) –id узла, который необходимо закрыть

 

connected(node_id)

 

Если подключен node_id, возвращаент значение True

 

connection_delay(node_id)

 

Возвращает количество миллисекунд ожидания (в зависимости от состояния подключения) перед отправкой данных. При отключении учитывается время задержки повторного подключения. При подключении возвращает 0 для того, чтобы разрешить завершение неблокирующего подключения. При подключении возвращает очень большое число для обработки медленных/остановленных подключений.

Параметры:

node_id (int) – id узла, который необходимо проверить

Выводы:

Время ожидания в миллисикундах

Тип вывода:

int

 

get_api_versions()

Возвращает карту ApiVersions, если она доступна.

Примечание: Вызов check_version должен быть выполнен, необходимо использовать версию 0.10.0 или более новую

Возвращает: карту отображения dict {api_key: (min_version, max_version)} или None (если ApiVersion не поддерживается кластером kafka).

 

in_flight_request_count(node_id=None)

Выводит  количество запросов для 1 узла или всех узлов.

Параметры:

node_id (int, optional) – узел, который необходимо  проверить. Если не указано иное, возвращает итоговое значение, которое является общим для всех узлов

Выводы:

запросы для 1 узла или всех узлов

Тип вывода:

int

 

is_disconnected(node_id)

Проверяет, отсоединено ли соединение узла или нет.

Отсоединенный узел либо был закрыт, либо зафиксорован сбой. Отказы соединения обычно являются кратковременными и могут быть возобновлены при следующем вызове метода ready (), но бывают случаи, когда кратковременные отказы необходимо выявить и устранить повторно.

Параметры:

node_id (int) –id узла, который необходимо проверить

Выводы:

Если узел существует и отклюен, возвращается значение True

Тип вывода:

bool

 

is_ready(node_id, metadata_priority=True)

 

Проверьте, готов ли к отправке дополнительных запросов выбранный узел.

В дополнение к проверкам на уровне подключения этот метод также используется для блокировки отправки дополнительных запросов во время обновления метаданных.

Параметры:

  • node_id (int) – id узла, который необходимо проверить
  • metadata_priority (bool) – Помечает узел как не готовый к работе, если требуется обновление метаданных. По умолчанию: True

Выводы:

True, если узел готов к работе,  а метаданные не были обновлены

Тип вывода:

bool

 

least_loaded_node()

 

Выберите узел с наименьшим количеством невыполненных запросов и резервными версиями.

Этот метод предпочтет узел с существующим соединением и без запросов. Если такой узел не найден, узел будет выбран случайным образом из разъединенных узлов, которые не «затемнены» (т.е. не подлежат откату при повторном подключении). Если метаданные узла не были получены, будет возвращен узел начальной загрузки (при экспоненциальном откате).

Выводы:

node_id или None, если не удалось найти подходящий узел

 

maybe_connect(node_id, wakeup=True)

Постановка узла в очередь для асинхронного соединения во время следующего Poll ()

 

poll(timeout_ms=None, future=None)

 

Пробует читать и записывать в сокеты.

Этот метод также попытается завершить подключения узлов, обновить устаревшие метаданные и выполнить ранее запланированные задачи.

Параметры:

  • timeout_ms (int, optional) максимальное время ожидания (в мс) хотя бы одного ответа. Должно быть неотрицательным. Фактическое время ожидания будет минимальным из времени ожидания, времени ожидания запроса и времени ожидания метаданных. По умолчанию: request_timeout_ms
  • future (Future, optional) – если предусмотрено, блокируется до future.is_done

Выводы:

Ответ получен (может быьб пустым)

Тип вывода:

list

 

ready(node_id, metadata_priority=True)

 

Проверяет, подключен ли узел и можно ли отправлять дополнительные запросы.

Параметры:

  • node_id (int) - id узла, оторый необходимо проверить
  • metadata_priority (bool) – Помечает узел как не готовый к работе, если требуется обновление метаданных. По умолчанию: True

Выводы:

True, если мы готовы отправить данные на этот узел

Тип вывода:

bool

 

send(node_id, request, wakeup=True)

 

Отправляет запрос на определенный узел. Байты помещаются во внутреннюю очередь отправки для каждого подключения. Фактический сетевой ввод-вывод будет инициирован при последующем вызове .poll ()to .poll()

Параметры:

  • node_id (int) – узел назначения
  • request (Struct) – объект запроса (незакодированный)
  • wakeup (bool) – флаг для отключения thread-wakeup

Указывает:

AssertionError – если node_id отсутствует в текущих метаданных кластера

Выводы:

устраняет структуру ответа или ошибку

Тип вывода:

Future

 

set_topics(topics)

 

Задает конкретные топики для отслеживания метаданных.

Параметры:

topics (list of str) топики для проверки метаданных

Выводы:

появляется после запроса/ответа метаданных

Тип вывода

Future

 

BrokerConnection

classkafka.BrokerConnection(host, port, afi, **configs)

Инициализирует соединение с брокером Kafka

Ключевые аргументы:

 

  • client_id (str) – имя этого клиента. Эта строка передается в каждом запросе к серверам и может использоваться для идентификации определенных записей журнала на стороне сервера, которые соответствуют этому клиенту. Также представлен координатору группы для регистрации в отношении администрирования групп консюмеров. По умолчанию: ‘kafka-python-{version}’
  • reconnect_backoff_ms (int) – Время ожидания повторного подключением к данному узлу в миллисекундах. По умолчанию: 50.
  • reconnect_backoff_max_ms (int) – Максимальное время (в миллисекундах) задержки/ожидания при повторном подключении к брокеру, которому не удалось подключиться ранеее. Если этот параметр установлен, то при каждом последовательном сбое соединения величина задержки будет увеличиваться до достижения максимума. После достижения максимального значения попытки повторного подключения будут продолжаться с определенной периодичностью. Чтобы избежать штормов соединения, к задержке будет применяться коэффициент рандомизации 0,2, что приведет к диапазону: от 20% ниже до 20% выше расчетного значения. По умолчанию: 1000.
  • request_timeout_ms (int) – Время ожидания запроса клиента в миллисекундах. По умолчанию:30000.
  • max_in_flight_requests_per_connection (int) – Запросы передаются брокерам kafka до этого количества максимальных запросов на соединение с брокером. По умолчанию:5.
  • receive_buffer_bytes (int) - Размер буфера приема TCP (SO_RCVBUF), используемого при чтении данных. По умолчанию: None (зависит от системных настроек по умолчанию). Клиент Java по умолчанию имеет значение 32768.
  • send_buffer_bytes (int) – Размер буфера передачи TCP (SO_SNDBUF), используемого при передаче данных. По умолчанию: None (зависит от системных настроек по умолчанию). Клиент Java по умолчанию имеет значение 131072.
  • socket_options (list) – Список аргументов кортежа для socket.setsockopt, применяемых к сокетам подключений брокера. По умолчанию: [(socket.IPPROTO_TCP, сокет. TCP_NODELAY, 1)]
  • security_protocol (str) – Протокол, используемый для связи с брокерами. Допустимые значения: PLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL. По умолчанию: PLAINTEXT.
  • ssl_context (ssl.SSLContext) - предварительно настроенный SSLContext для обертывания соединений сокетов. Если указано, все остальные конфигурации ssl _ * будут игнорироваться. По умолчанию: None.
  • ssl_check_hostname (bool) – флаг, чтобы настроить, должно ли ssl handshake проверять, соответствует ли сертификат имени хоста брокеров. По умолчанию: True.
  • ssl_cafile (str) – Наименование файла ca, которое будет испольоваться при проверке сертификата. По умолчанию: None.
  • ssl_certfile (str) – необязательное имя файла в формате pem, содержащего сертификат клиента, а также любые сертификаты ca, необходимые для установления подлинности сертификата. по умолчанию: Нет.
  • ssl_keyfile (str) – необязательное имя файла, содержащее закрытый ключ клиента. По умолчанию: None.
  • ssl_password (callable, str, bytes, bytearray) – необязательный пароль или вызываемая функция, возвращающая пароль, для расшифровки закрытого ключа клиента. По умолчанию: None.
  • ssl_crlfile (str) – необязательное имя файла, содержащее CRL для проверки срока действия сертификата. По умолчанию проверка CRL не выполняется. При предоставлении файла по этому CRL будет проверяться только конечный сертификат. CRL можно проверить только с помощью Python 3.4 + или 2.7.9 +. По умолчанию: Нет.
  • ssl_ciphers (str) – опционально устанавливать доступные шифры для ssl-соединений. Это должна быть строка в формате списка шифров OpenSSL. Если ни один шифр не может быть выбран (поскольку параметры compile-time или другая конфигурация запрещают использование всех указанных шифров), ssl. Будет поднят SSLError. Смотрите ssl. SSLContext.set_ciphers
  • api_version (tuple) – Укажите, какую версию API Kafka использовать. Допустимые значения: (0, 8, 0), (0, 8, 1), (0, 8, 2), (0, 9), (0, 10). По умолчанию: (0, 8, 2)
  • api_version_auto_timeout_ms (int) – количество миллисекунд, необходимое для выдачи исключения тайм-аута из конструктора при проверке версии API брокера. Применяется только в том случае, если api_version имеет значение None
  • selector (selectors.BaseSelector) – Предоставьте конкретную реализацию селектора для мультиплексирования ввода-вывода. По умолчанию: selectors.DefaultSelector
  • state_change_callback (callable) – функция, вызываемая при изменении состояния соединения с CONNECTING на CONNECTED и т.д..
  • metrics (kafka.metrics.Metrics) – Дополнительно предоставьте экземпляр метрики для сбора статистики ввода-вывода сети. По умолчанию: None.
  • metric_group_prefix (str) – префикс для названий метрик. По умолчанию: ‘’
  • sasl_mechanism (str) – Механизм аутентификации, когда security_protocol настроен для SASL_PLAINTEXT или SASL_SSL. Допустимые значения: PLAIN, GSSAPI, OAUTHBEARER, SCRAM-SHA-256, SCRAM-SHA-512.
  • sasl_plain_username (str) – имя пользователя для аутентификации sasl PLAIN и SCRAM. Требуется, если sasl_mechanism является PLAIN или одним из механизмов SCRAM.
  • sasl_plain_password (str) – пароль для аутентификации sasl PLAIN и SCRAM. Требуется, если sasl_mechanism является PLAIN или одним из механизмов SCRAM.
  • sasl_kerberos_service_name (str) – Имя службы, включаемое в квитирование механизма sasl GSSAPI. По умолчанию: kafka
  • sasl_kerberos_domain_name (str) – доменное имя kerberos для использования при квитировании механизма GSSAPI sasl. По умолчанию: один из серверов начальной загрузки
  • sasl_oauth_token_provider (AbstractTokenProvider) – Экземпляр поставщика маркеров OAuthBearer. (См. Kafka.oauth.abstract). По умолчанию: None

 

blacked_out()

 

Если мы отключились от данного узла и пока не можем восстановить соединение,  возвращает true.

 

can_send_more()

 

Если нет max_in_flight_requests_per_connection , возвращает True.

 

check_version(timeout=2, strict=False, topics=[])

 

Делает попытку угадать версию брокера.

Примечание: Это блокирующий вызов.

Возвращает: кортеж версии, т.е. (0, 10), (0, 9), (0, 8, 2),...

 

close(error=None)

 

Закрывает сокет и не выполняет

Параметры:

error (Exception, optional) – ожидающие запросы будут отклонены. По умолчанию: kafka.errors.KafkaConnectionError.

 

connect()

Попытка подключения и возврата ConnectionState

 

connected()

 

Возвращает True в случае, если сокет подключен.

 

connecting()

 

Возвращает True, если соединение все еще выполняется (процесс подключения может включать в себя несколько различных состояний, таких как подтверждение SSL, авторизация и т.д.).

 

connection_delay()

 

Возвращет количество миллисекунд ожидания (в зависимости от состояния подключения) перед отправкой данных. При отключении учитывается время задержки повторного подключения. При подключении возвращает очень большое число для обработки медленных/остановленных подключений.

 

disconnected()

 

Если сокет закрыт, возвращает True.

 

recv()

 

Получение неблокирующей сети.

Возвращаемый список кортежей (response, future)

 

send(request, blocking=True)

 

Запрос очереди на асинхронную сеть

 

send_pending_requests()

 

Попытки отправки сообщений о ожидающих запросах через блокирование ввода-вывода. Если все запросы были успешно отправлены, возвращает True. В противном случае, если сокет заблокирован и есть еще байты для отправки, возвращает False.

 

send_pending_requests_v2()

 

Попытки отправки сообщений об ожидающих запросах через неблокирующий IO. Если все запросы были успешно отправлены, возвращает True. В противном случае, если сокетзаблокирован и есть еще байты для отправки, возвращает False.

 

ClusterMetadata

classkafka.cluster.ClusterMetadata(**configs)

Класс для управления метаданными кластера Kafkа.

Этот класс не выполняет никаких операций ввода-вывода, он просто обновляет внутреннее состояние, заданное ответами API (MetadataResponse, GroupCoordinatorResponse).

Ключевые аргументы:

 

  • retry_backoff_ms (int) – Миллисекунды для отката при повторении ошибок. По умолчанию: 100.
  • metadata_max_age_ms (int) – Период времени в миллисекундах, после которого мы принудительно обновляем метаданные, даже если мы не видели никаких изменений в лидерстве разделов, чтобы проактивно обнаруживать какие-либо новые брокеры или разделы. По умолчанию: 300000
  • bootstrap_servers – ‘host[:port]’ строка (или список строк «host [: port]»), к которой клиент должен обратиться для загрузки начальных метаданных кластера. Это не обязательно должен быть полный список узлов. Ему просто нужно иметь хотя бы одного брокера, который будет отвечать на запрос Metadata API. Порт по умолчанию - 9092. Если серверы не указаны, по умолчанию будет установлено значение localhost: 9092.

 

add_group_coordinator(group, response)

 

Обновление метаданных для координатора группы

Параметры:

  • group (str) – название группы из GroupCoordinatorRequest
  • response (GroupCoordinatorResponse) – ответ брокера

Выводы:

Если метаданные обновлены, координатор node_id. None или ошибка

Тип вывода:

строка

 

add_listener(listener)

 

Добавление функции обратного вызова, вызываемой при каждом обновлении метаданных

 

available_partitions_for_topic(topic)

 

Возвращает набор партиций с известными лидерами

Параметры:

topic (str) – топик, который необходимо проверить на наличие разделов

Выводы:

{partition (int), …} Если тпик не найден, то None

Тип вывода:

set

 

broker_metadata(broker_id)

 

Получает BrokerMetadata

Параметры:

broker_id (int) – node_id для брокера, которого необходимо проверить

Выводы:

BrokerMetadata или None

 

brokers()

 

Получает BrokerMetadata

Выводы:

{BrokerMetadata, …}

Тип вывода:

set

 

coordinator_for_group(group)

 

Возвращает node_id координатора группы

Параметры:

group (str) – название группы консюмеров

Выводы:

node_id для координатора группы, если группа не существует, то None

Тип вывода:

int

 

failed_update(exception)

 

Обновляет состояния кластера при сбое запроса MetadataRequest.

 

leader_for_partition(partition)

 

Возврат node_id лидера, -1 недоступн и None, если значение неизвестно.

 

partitions_for_broker(broker_id)

 

Вернуть TopicPartitions, для которых брокер является лидером.

Параметры:

broker_id (int) –id узла для брокера

Выводы:

{TopicPartition, …} Отсутствует, если брокер не имеет разделов или не существует.

Тип вывода:

set

 

partitions_for_topic(topic)

 

Возвращает набор всех разделов для топика (доступно или нет)

Параметры:

topic (str) – топик, который нужно проверить на наличие разделов

Выводы:

{partition (int), …}

Тип вывода:

set

 

refresh_backoff()

 

Возвращает миллисекунды для ожидания перед попыткой повторить попытку после сбоя

 

remove_listener(listener)

 

Удаляет ранее добавленный обратный вызов прослушивателя

 

request_update()

 

Флаги метаданных для обновления, возвращает Future ()

Фактическое обновление должно обрабатываться отдельно. Этот метод изменит только сообщаемое значение ttl ()

Выводы:

kafka.future.Future (значением будет объект кластера после обновления)

 

topics(exclude_internal_topics=True)

 

Выводит набор известных топиков.

Параметры:

exclude_internal_topics (bool) – определяет, должны ли записи из внутренних топиков (таких как смещения) быть доступны консюмеру. Если установлено значение True, единственным способом получения записей из внутреннего раздела является подписка на него. По умолчанию: True

Выводы:

{topic (str), …}

Тип вывода:

set

 

ttl()

 

Миллисекунды до обновления метаданных

 

update_metadata(metadata)

 

Обновляет состояние кластера для MetadataResponse.

Параметры:

metadata (MetadataResponse) – ответ брокера на запрос метаданных

 

with_partitions(partitions_to_add)

 

Возвращает копию кластера метаданных с включенными в него разделами

 

Установка

Вы можете установить kafka-python с помощью любого предпочитаемого Вами менеджера.

Новая версия

Pip:

pip install kafka-python

 

Все доступные версии перечислены здесь: https://github.com/dpkp/kafka-python/releases

 

Самая последняя версия

git clone https://github.com/dpkp/kafka-python
pip install ./kafka-python

 

Дополнительная установка LZ4

Чтобы запустить сжатие/распаковку LZ4, установите python-lz4:

>>> pip install lz4

 

Дополнительная установка crc32c

Чтобы включить оптимизированную проверку контрольной суммы CRC32, установите crc32c:

>>> pip install crc32c

 

Дополнительная установка Snappy

Загрузите  Snappy из https://google.github.io/snappy/

Ubuntu:

apt-get install libsnappy-dev

OSX:

brew install snappy

 

Из источника:

wget https://github.com/google/snappy/releases/download/1.1.3/snappy-1.1.3.tar.gz
tar xzvf snappy-1.1.3.tar.gz
cd snappy-1.1.3
./configure
make
sudo make install

 

Установка модуля Python

Установите модуль python-snappy:

pip install python-snappy

 

Дополнительная установка  crc32c

Настоятельно рекомендуется, если Вы используете брокеров Kafka 11 +. Для него kafka-python использует новую версию протокола сообщений, которая требует вычисления crc32c, что отличается от реализации хеша zlib.crc32. По умолчанию kafka-python вычисляет его на чистом python, что довольно медленно. Чтобы ускорить его, мы опционально выбираем следующий пакет: https://pypi.python.org/pypi/crc32c.

pip install crc32c

 

Тесты

Управление тестовыми средами осуществляется через tox. Тестовый набор запускается через pytest.

Линтинг выполняется через pylint, но обычно пропускается из-за проблем совместимости/производительности pylint.

Для получения подробной информации см. здесь: https://coveralls.io/github/dpkp/kafka-python

Набор тестов включает в себя модульные тесты, которые имитируют сетевые интерфейсы, а также интеграционные тесты, которые настраивают и снимают устройства брокера kafka (и zookeeper) для тестирования клиента/консюмера/продюсера.

 

Модульные тесты

Для запуска этих тестов установите tox:

pip install tox

 

Подробная информация: https://tox.readthedocs.io/en/latest/install.html

Затем просто запустите tox, установив среду python.

tox -e py27
tox -e py35
 
# run protocol tests only
tox -- -v test.test_protocol
 
# re-run the last failing test, dropping into pdb
tox -e py27 -- --lf --pdb
 
# see available (pytest) options
tox -e py27 -- --help

 

Интеграционные тесты

KAFKA_VERSION=0.8.2.2 tox -e py27
KAFKA_VERSION=1.0.1 tox -e py36

 

Интеграционные тесты запускают Kafka и Zookeeper. Для этого также требуется загрузка серверных двоичных файлов kafka:

./build_integration.sh

 

По умолчанию это установит версии брокера, перечисленные в build_integration.sh's ALL_RELEASES, в сервер/каталог. Чтобы установить определенную версию, установите переменную KAFKA_VERSION:

KAFKA_VERSION=1.0.1 ./build_integration.sh

 

Затем, чтобы запустить тесты для конкретной версии Kafka, установите переменную KAFKA_VERSION env в сборку сервера, которую Вы хотите использовать для тестирования:

KAFKA_VERSION=1.0.1 tox -e py36

 

Для тестирования исходного дерева kafka установите KAFKA_VERSION=trunk [дополнительно установите SCALA_VERSION (по умолчанию значение, установленное в build_integration.sh)]

SCALA_VERSION=2.12 KAFKA_VERSION=trunk ./build_integration.sh
KAFKA_VERSION=trunk tox -e py36

 

Совместимость

kafka-python совместим с брокером версий с 2.4 по 0.8.0. kafka-python не совместим с версией0.8.2-beta.

Поскольку протокол сервера kafka является обратно совместимым, ожидается, что kafka-python будет работать и с более новыми версиями брокера.

Хотя kafka-python протестирован и, как ожидается, будет работать на самых последних версиях брокера, все же не все функции будут поддерживаться. В частности, кодеки аутентификации и транзакционная поддержка продюсеров/консюмеров будут реализованы не полностью. Пиарщики приветствуются =)

kafka-python тестируется на python 2.7, 3.4, 3.7 и pypy2.7.

Тестирование через Travis-CI: https://travis-ci.org/dpkp/kafka-python

 

Полезные ссылки

Информация на Github: https://github.com/dpkp/kafka-python

Ограниченный IRC чат на # kafka-python на freenode (общий чат # apache-kafka).

Общая информация касательно Apache Kafka доступна по ссылке: https://kafka.apache.org/

Обсуждение проекта и реализации kafka-клиента (не специфичного для python): https://groups.google.com/forum/m/#!forum/kafka-clients

 

 

Узнать стоимость решенияЗапросить видео презентацию

← Предыдущая статья
Начало работы с Apache Kafka в Docker: пошаговая инструкция
Следующая статья →
Python-клиент для Apache Kafka
Запросить видео презентацию Запросить доступ к демо стенду online Узнать стоимость лицензий

Задать вопрос

loading...

Решения

Анализировать ФинансыУвеличивайте ПродажиОптимальный Склад и ЛогистикаМаркетинговые Метрики

Клиенты
  • KERAMA MARAZZI — международный бренд, входящий в число лидеров глобального рынка керамики. Бизнес компании охватывает весь процесс создания керамических изделий, от глиняных карьеров до фирменной розницы во всех крупных городах РФ и за рубежом.

  • АО «Евросиб СПб–транспортные системы» – оператор контейнерных сервисов с широкой сетью маршрутов на внутрироссийских и международных направлениях. Имеет успешный опыт управления парком фитинговых платформ, а также организации ускоренных контейнерных поездов, в основе которых точное расписание, оптимальные сроки доставки груза и экономическая целесообразность.

  • "Холодильник.ру" - крупнейший в России интернет-магазин бытовой техники и электроники. Компания была основана в 2003 году и за почти 20 лет работы завоевала лидирующие позиции на рынке онлайн ритейла. По данным исследовательского агентства Data Insight, "Холодильник.ру" входит в top-10 крупнейших интернет-магазинов России в категории "электроника и бытовая техника". Компания имеет развитую логистическую инфраструктуру и ежедневно осуществляет более 3500 доставок заказов по всей стране.

  • ГК «Агропромкомплектация-Курск» - одна из ведущих в Российской Федерации агропромышленных компаний с полным производственным циклом "от поля до прилавка". За 32 года работы на рынке компания заслуженно завоевала репутацию одного из лидеров страны в производстве свинины и молока.

  • Решения
    • Дистрибуция
    • Розничная торговля
    • Производство
    • Операторы связи
    • Страхование
    • Банки
    • Лизинг
    • Логистика
    • Нефтегазовый сектор
    • Медицина
    • Сеть ресторанов
    • E-Commerce
    • Энергетика
    • Фармацевтика
  • Услуги
    • Переход на отечественные BI и DWH
    • Консалтинг
    • Пилотный проект
    • Обучение и сертификация
    • Бесплатное обучение
    • Техническая поддержка
    • Технические задания
    • Сбор требований для проекта внедрения BI-системы
    • CI/CD для DWH
    • Аудит BI приложений
    • Выделенная команда
    • Настойка и поддержка баз данных
    • Разработка BI Стратегии
    • Styleguide для BI-системы
    • Как выбрать BI-систему
  • Платформы
    • FineBI
    • FineReport
    • FineDataLink
    • Коннекторы данных из 1С в BI
    • Airflow + NiFi
    • Visiology
    • Luxms BI
    • Modus BI
    • PIX BI
    • Arenadata
    • ClickHouse
    • Greenplum
    • Postgres Professional
    • Open-source BI: Superset/Metabase
    • Loginom
    • Yandex.DataLens
    • AI / Исскуственный интеллект
    • Optimacros
    • Шины данных
  • Курсы
    • Учебный курс Информационная грамотность
    • Учебный курс для бизнес-аналитиков
    • Учебный курс для системных аналитиков
    • Учебный курс по Data Governance
    • Учебный курс Как стать CDO
    • Учебный курс Современная архитектура хранилища данных
    • Учебный курс по Fine BI
    • Учебный курс по FineReport
    • Учебный курс по DWH
    • Учебный курс по Data Science (ML, AI)
    • Учебный курс по PostgreSQL
    • Учебный курс по Apache Airflow и NiFi
    • Учебный курс по Open-source BI
    • Учебный курс по ClickHouse
    • Учебный курс по DataLens
    • Учебный курс по Loginom
    • Учебный курс по Modus BI и ETL
    • Учебный курс по Visiology
    • Учебный курс по dbt
  • Функциональные решения
    • Создание Data Lake
    • Цифровая трансформация
    • Управление по KPI
    • Финансы
    • Продажи
    • Склад
    • HR
    • Маркетинг
    • Внутренний аудит
    • Категорийный менеджмент
    • S&OP и прогнозная аналитика
    • Геоаналитика
    • Цепочки поставок (SCM)
    • AutoML
    • Process Mining
    • Сквозная аналитика
  • Компания
    • О нас
    • Руководство
    • Новости
    • Клиенты
    • Скачать
    • Контакты
    • Политика конфиденциальности
RutubeVkontakteLinkedInYouTube
ООО "Би Ай Консалт",
ИНН: 7811437757,
ОГРН: 1097847154184
199178, Россия,
Санкт-Петербург,
6-ая линия В.О., Д. 63, 4 этаж
Тел: +7 (812) 334-08-01
Тел: +7 (499) 608-13-06
E-mail: info@biconsult.ru

 

 

 

 

 

×

Пользуясь сайтом, вы соглашаетесь с использованием cookies и политикой конфиденциальности.