kafka-python
Введение
kafka-python разработан таким образом, чтобы быть в состоянии функционировать подобно официальному java-клиенту, только с добавлением питоновских интерфейсов (например, потребительских итераторов).
kafka-python лучше всего использовать с новыми брокерами (0.9+), но он совместим и с более старыми версиями (до 0.8.0). Некоторые функции будут доступны только в новых брокерах. Например, полностью скоординированные группы потребителей - т.е. динамическое назначение разделов нескольким потребителям в одной группе - требуют использования брокеров kafka версии 0.9. Применение этой функции в более ранних версиях брокера потребует написания специального кода (возможно, с использованием zookeeper или consul). Для более старых брокеров Вы можете добиться чего-то подобного, вручную назначив различные разделы каждому экземпляру потребителя с помощью инструментов управления конфигурацией, таких как chef, ansible и т. д. Такой подход сработает несмотря на то, что он не поддерживает ребалансировку при сбоях. Более подробную информацию см. в разделе «Совместимость».
Обратите внимание, что мастер-ветка может содержать невыпущенные функции. Документацию по релизу можно найти в readthedocs и/или встроенной справке python.
>>> pip install 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() перед непосредственным потреблением записей. |
|
Ключевые аргументы |
Примеры
По умолчанию: None
|
ПРИМЕЧАНИЕ:
Подробное описание параметров настройки доступно в https://kafka.apache.org/documentation/#consumerconfigs
assign(partitions)
Присвоение TopicPartitions консюмеров вручную.
|
Параметры: |
partitions (list of TopicPartition) – назначение для данного экземпляра |
|
Указывает: |
|
ВНИМАНИЕ!
Невозможно одновременно использовать назначение партиций с помощью 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.
Это асинхронный вызов, поэтому блокироваться он не будет. Любые возникшие ошибки либо передаются обратному вызову (если он предусмотрен), либо отбрасываются:
|
Параметры: |
|
|
Выводы: |
kafka.future.Future |
committed(partition, metadata=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 |
|
Указывает: |
|
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 будет произвольно использован в середине потребления для сброса смещений выборки.
|
Параметры: |
|
|
Указывают: |
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().
|
Параметры: |
В рамках управления группами консюмер отслеживает список консюмеров, принадлежащих к определенной группе, и запускает операцию ребалансировки, как толькр происходит одно из следующих событий:
При наступлении любого из этих событий предоставленный слушатель будет вызван для указания того, что назначение консюмера было отозвано. При этом есть гарантия того, что разделы, отозванные/назначенные через этот интерфейс, относятся к топикам, подписанным в этом вызове. |
|
Указывает: |
|
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 указывают, как превратить объекты ключа и значения, предоставленные пользователем, в байты.
|
Ключевые аргумерты: |
|
|
|
Количество подтверждений, которое прод.сер должен получить от лидера, прежде чем считать запрос завершенным. Это позволяет контролировать долговечность отправляемых записей. Обычно используются следующие настройки:
0: прод.сер не будет ждать подтверждений от сервера. Сообщение будет немедленно добавлено в буфер сокета и будет считаться отправленным. В этом случае нет гарантии, что сервер получил запись, конфигурация повторных попыток не будет действовать (поскольку клиент обычно не знает о неудачах). Смещение, возвращаемое для каждой записи, всегда будет равно -1.
1: Дожидается, пока лидер не запишет запись в свой локальный журнал. Брокер ответит, не дожидаясь полного подтверждения от всех последователей. В этом случае, если лидер выйдет из строя сразу после подтверждения записи, но до того, как ее воспроизведут последователи, запись будет потеряна. all: Дожидается того момента, когда полный набор синхронизированных реплик запишет запись. Это гарантирует, что запись не будет потеряна до тех пор, пока «жива» хотя бы одна синхронизированная реплика. Это самая надежная из всех доступных гарантий. Если значение не установлено, по умолчанию используется acks=1.
|
ПРИМЕЧАНИЕ:
Более подробная информация качательно параметров настройки доступна по ссылке: 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)
Публикует сообщение в топик.
|
Параметры |
|
|
Выводы: |
Относится к 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)
Создание дополнительных разделов для существующего топика.
|
Параметры: |
|
|
Выводы: |
Необходимая версия класса CreatePartitionsResponse. |
create_topics(new_topics, timeout_ms=None, validate_only=False)
Создает новый топик в кластере.
|
Параметры: |
|
|
Выводы: |
Необходимая версия класса CreateTopicResponse |
delete_acls(acl_filters)
Удаляет набор ACL
Удаляет все ACL, соответствующие списку входных ACLFilter
|
Параметры: |
|
|
Выводы: |
Список, состоящий из 3 кортежей, соответствующих списку входных фильтров. Кортежи содержат (входной ACLFilter, список затронутых ACL, а также экземпляр KafkaError) |
Удаляет топики из кластера.
|
Параметры: |
|
|
Выводы: |
Необходимая версия класса DeleteTopicsResponse. |
describe_acls(acl_filter)
Описывает набо ACL
Используется для возврата набора ACL, соответствующих заданному фильтру ACLFilter. Для этого кластер должен быть сконфигурирован с авторизатором, иначе Вы получите ошибку SecurityDisabledError
|
Параметры: |
|
|
Возвраты: |
кортеж из списка подходящих объектов ACL и ошибки KafkaError (NoError в случае успеха) |
describe_configs(config_resources, include_synonyms=False)
Получение параметров конфигурации для одного или нескольких ресурсов Kafka.
|
Параметры: |
|
|
Выводы: |
Необходимая версия DescribeConfigsResponse. |
describe_consumer_groups(group_ids, group_coordinator_id=None, include_authorized_operations=False)
Описывает набор групп консюмеров.
Сообщения о любых ошибках приходят немедленно.
|
Параметры: |
|
|
Выводы: |
Список описаний групп. На данный момент описания групп - это необработанные результаты из ответа DescribeGroupsResponse. В перспективе мы планируем изменить его для того, чтобы он возвращал именованные кортежи, а также декодировал назначения разделов. |
list_consumer_group_offsets(group_id, group_coordinator_id=None, partitions=None)
Получение смещений консюмеров для одной группы консюмеров.
Примечание: При этом не проверяется, действительно ли group_id существуют в кластере.
Как только возникает какая-либо ошибка, сразу же приходит соответствующее сообщение.
|
Параметры: |
|
|
Выводы: |
Словарь с ключами 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».
Как только возникает какая-либо ошибка, сразу же приходит соответствующее сообщение.
|
Параметры: |
|
|
Выводы: |
Список кортежей групп консюмеров. |
|
Указывает: |
|
KafkaClient
classkafka.KafkaClient(**configs)
Сетевой клиент для асинхронных запросов/ответов при вводе/выводе данных по сети.
Это внутренний класс, используемый для реализации пользовательских клиентов продюсера и консюмера.
Этот класс не является потокобезопасным!
cluster
Локальный кэш метаданных кластера, получаемый через MetadataRequests во время poll().
|
Тип: |
ClusterMetadata |
|
Ключевые аргументы: |
|
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),...
|
Указывает: |
|
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)
Проверьте, готов ли к отправке дополнительных запросов выбранный узел.
В дополнение к проверкам на уровне подключения этот метод также используется для блокировки отправки дополнительных запросов во время обновления метаданных.
|
Параметры: |
|
|
Выводы: |
True, если узел готов к работе, а метаданные не были обновлены |
|
Тип вывода: |
bool |
least_loaded_node()
Выберите узел с наименьшим количеством невыполненных запросов и резервными версиями.
Этот метод предпочтет узел с существующим соединением и без запросов. Если такой узел не найден, узел будет выбран случайным образом из разъединенных узлов, которые не «затемнены» (т.е. не подлежат откату при повторном подключении). Если метаданные узла не были получены, будет возвращен узел начальной загрузки (при экспоненциальном откате).
|
Выводы: |
node_id или None, если не удалось найти подходящий узел |
maybe_connect(node_id, wakeup=True)
Постановка узла в очередь для асинхронного соединения во время следующего Poll ()
poll(timeout_ms=None, future=None)
Пробует читать и записывать в сокеты.
Этот метод также попытается завершить подключения узлов, обновить устаревшие метаданные и выполнить ранее запланированные задачи.
|
Параметры: |
|
|
Выводы: |
Ответ получен (может быьб пустым) |
|
Тип вывода: |
list |
ready(node_id, metadata_priority=True)
Проверяет, подключен ли узел и можно ли отправлять дополнительные запросы.
|
Параметры: |
|
|
Выводы: |
True, если мы готовы отправить данные на этот узел |
|
Тип вывода: |
bool |
send(node_id, request, wakeup=True)
Отправляет запрос на определенный узел. Байты помещаются во внутреннюю очередь отправки для каждого подключения. Фактический сетевой ввод-вывод будет инициирован при последующем вызове .poll ()to .poll()
|
Параметры: |
|
|
Указывает: |
AssertionError – если node_id отсутствует в текущих метаданных кластера |
|
Выводы: |
устраняет структуру ответа или ошибку |
|
Тип вывода: |
Future |
set_topics(topics)
Задает конкретные топики для отслеживания метаданных.
|
Параметры: |
topics (list of str) топики для проверки метаданных |
|
Выводы: |
появляется после запроса/ответа метаданных |
|
Тип вывода |
Future |
BrokerConnection
classkafka.BrokerConnection(host, port, afi, **configs)
Инициализирует соединение с брокером Kafka
|
Ключевые аргументы: |
|
|
|
|
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).
|
Ключевые аргументы: |
|
|
|
|
add_group_coordinator(group, response)
Обновление метаданных для координатора группы
|
Параметры: |
|
|
Выводы: |
Если метаданные обновлены, координатор 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




