Python-клиент для Apache Kafka
Confluent, ведущий разработчик Apache Kafka, предлагает пользователям confluent-kafka-python на GitHub. Этот Python-клиент предоставляет высокоуровневых продюсеров, консюмеров и AdminClient, совместимых с брокерами Kafka (версия 0.8 и новее), Confluent Cloud и Confluent Platform. Следите за последними изменениями, проверяя список изменений, доступный в том же репозитории.
Установка Python-клиента
Клиент доступен на PyPI и может быть установлен с помощью pip:
pip install confluent-kafka
Вы можете установить его в виртуальном приложении virtualenv .
ПРИМЕЧАНИЕ:
Пакет confluent-kafka-python поставляется в комплекте с предварительно установленной версией librdkafka, которая не включает поддержку GSSAPI/Kerberos.
Продюсер Kafka
Инициализация
Продюсер настраивается согласно инструкции, указанной ниже.
Если Вы запускаете Kafka локально, Вы можете инициализировать продюсера следующим образом:
from confluent_kafka import Producer
import socket
conf = {'bootstrap.servers': 'host1:9092,host2:9092',
'client.id': socket.gethostname()}
producer = Producer(conf)
Если Вы подключаетесь к кластеру Kafka в Confluent Cloud, Вам необходимо предоставить учетные данные. В приведенном ниже примере показано использование ключа и секрета API кластера.
from confluent_kafka import Producer
import socket
conf = {'bootstrap.servers': 'pkc-abcd85.us-west-2.aws.confluent.cloud:9092',
'security.protocol': 'SASL_SSL',
'sasl.mechanism': 'PLAIN',
'sasl.username': '<CLUSTER_API_KEY>',
'sasl.password': '<CLUSTER_API_SECRET>',
'client.id': socket.gethostname()}
producer = Producer(conf)
Асинхронная запись
Для инициирования отправки сообщения в Kafka вызовите метод product, передав значение сообщения (которое может быть равно None), ключ (необязательно), раздел и обратный вызов. Вызов Product завершается немедленно. Если сообщение не может быть помещено в очередь из-за переполнения локальной очереди производства librdkafka, будет сгенерировано исключение KafkaException.
producer.produce(topic, key="key", value="value")
Для получения уведомления об успешной или неуспешной доставке можно передать параметр обратного вызова. Это может быть любой вызываемый объект, например, функция, связанный метод или вызываемый объект. Хотя метод produce () немедленно ставит в очередь сообщение для пакетной обработки, сжатия и передачи брокеру, никакие уведомления о доставке не будут распространяться до вызова метода poll ().
def acked(err, msg):
if err is not None:
print("Failed to deliver message: %s: %s" % (str(msg), str(err)))
else:
print("Message produced: %s" % (str(msg)))
producer.produce(topic, key="key", value="value", callback=acked)
# Wait up to 1 second for events. Callbacks will be invoked during
# this method call if the message is acknowledged.
producer.poll(1)
Синхронная запись
Клиент Python предоставляет метод flush (), который можно использовать для синхронной записи. Это, как правило, плохая идея, поскольку она эффективно ограничивает пропускную способность для брокера в оба конца, но в некоторых случаях она может быть оправдана.
producer.produce(topic, key="key", value="value") producer.flush()
Как правило, перед закрытием продюсера следует вызвать flush (), что обеспечит доставку всех невыполненных/поставленных в очередь сообщений.
Консюмер Kafka
Инициализация
Консюмер настраивается в соответствии с инструкцией, представленной ниже. Если Вы запускаете Kafka локально, Вы можете инициализировать консюмера так, как показано ниже.
from confluent_kafka import Consumer
conf = {'bootstrap.servers': 'host1:9092,host2:9092',
'group.id': 'foo',
'auto.offset.reset': 'smallest'}
consumer = Consumer(conf)
Если Вы подключаетесь к кластеру Kafka в Confluent Cloud, Вам необходимо предоставить учетные данные. В приведенном ниже примере показано использование ключа и секрета API кластера.
from confluent_kafka import Consumer
conf = {'bootstrap.servers': 'pkc-abcd85.us-west-2.aws.confluent.cloud:9092',
'security.protocol': 'SASL_SSL',
'sasl.mechanism': 'PLAIN',
'sasl.username': '<CLUSTER_API_KEY>',
'sasl.password': '<CLUSTER_API_SECRET>',
'group.id': 'foo',
'auto.offset.reset': 'smallest'}
consumer = Consumer(conf)
Свойство group.id является обязательным и определяет, в какую группу консюмеров входит тот или иной консюмер. Свойство auto.offset.reset указывает смещение, с которого конмюмер должен начать чтение в случае отсутствия фиксированных смещений для раздела или если фиксированное смещение недопустимо (возможно, из-за усечения журнала).
Локальный пример ниже показывает enable.auto.commit, настроенный на false в консюмере. Значение по умолчанию - True.
from confluent_kafka import Consumer
conf = {'bootstrap.servers': 'host1:9092,host2:9092',
'group.id': 'foo',
'enable.auto.commit': 'false',
'auto.offset.reset': 'earliest'}
consumer = Consumer(conf)
Примеры кода для Python-клиента
Базовый цикл
Типичное консюмерское приложение Kafka сосредоточено вокруг цикла потребления, который неоднократно вызывает метод опроса для извлечения записей по одному, которые были эффективно предварительно извлечены консюмером. Прежде чем войти в цикл потребления, Вы, как правило, используете метод подписки для того, чтобы указать, из каких топиков следует извлечь то или иное значение:
running = True
def basic_consume_loop(consumer, topics):
try:
consumer.subscribe(topics)
while running:
msg = consumer.poll(timeout=1.0)
if msg is None: continue
if msg.error():
if msg.error().code() == KafkaError._PARTITION_EOF:
# End of partition event
sys.stderr.write('%% %s [%d] reached end at offset %d\n' %
(msg.topic(), msg.partition(), msg.offset()))
elif msg.error():
raise KafkaException(msg.error())
else:
msg_process(msg)
finally:
# Close down consumer to commit final offsets.
consumer.close()
def shutdown():
running = False
Таймаут опроса жестко зафиксирован на 1 секунду. Если до истечения этого времени не получено ни одной записи, Consumer.poll () вернет пустой набор записей.
Обратите внимание на то, что Вы всегда должны вызывать Consumer.close () после того, как закончите использовать консюмера. Это обеспечит закрытие активных сокетов и очистку внутреннего состояния , а также немедленно вызовет перебалансировку группы, которая гарантирует, что любые разделы, принадлежащие консюмеру, будут переназначены другому члену группы. Если он не закрыт должным образом, брокер запустит перебалансировку только после истечения времени ожидания сеанса.
Синхронные коммиты
Самый простой и надежный способ ручной фиксации смещений - установка асинхронного параметра на вызов метода Consumer.commit (). Этот метод также может принимать взаимоисключающие смещения параметров ключевого слова для явного перечисления смещений для каждого назначенного раздела темы и сообщения, которые будут фиксировать смещения относительно объекта Message, возвращаемого poll ().
def consume_loop(consumer, topics):
try:
consumer.subscribe(topics)
msg_count = 0
while running:
msg = consumer.poll(timeout=1.0)
if msg is None: continue
if msg.error():
if msg.error().code() == KafkaError._PARTITION_EOF:
# End of partition event
sys.stderr.write('%% %s [%d] reached end at offset %d\n' %
(msg.topic(), msg.partition(), msg.offset()))
elif msg.error():
raise KafkaException(msg.error())
else:
msg_process(msg)
msg_count += 1
if msg_count % MIN_COMMIT_COUNT == 0:
consumer.commit(asynchronous=False)
finally:
# Close down consumer to commit final offsets.
consumer.close()
В этом примере синхронная фиксация инициирует каждые MIN_COMMIT_COUNT сообщения. Специальный флаг определяет, является ли этот вызов асинхронным или нет. Вы также можете инициировать фиксацию по истечении тайм-аута для того, чтобы гарантировать, что фиксированная позиция обновляется регулярно.
Гарантии доставки
Гарантия по умолчанию
Ниже приведены настройки Python - клиента по умолчанию, которые приводят к тому, что гарантия доставки по умолчанию клиента Python имеет значение «None»:
'enable.auto.commit': 'true' 'enable.auto.offset.store': 'true'
Поскольку автоматические фиксации выполняются в фоновом потоке, эти настройки могут привести к смещению последнего сообщения, зафиксированного до того, как приложение завершит его обработку. Если приложение должно было завершить работу в аварийном режиме или выйти до завершения обработки, и смещение было автоматически зафиксировано, следующее воплощение потребительского приложения начнется со следующего сообщения. В этом случае Вы можете потерять данные.
Гарантия “хотя бы один раз”
Для обеспечения данной гарантии необходимы следующие настройки:
'enable.auto.commit': 'true' 'enable.auto.offset.store': 'false'
Чтобы избежать потери данных или дубликатов с указанным выше режимом «None guaranteed» по умолчанию, приложение может отключить автоматическое хранилище смещений и вручную хранить смещения (с помощью rd_kafka_offsets_store ()) после обработки. Это дает приложению детальный контроль над тем, когда и как сообщение было зафиксировано. Последнее сохраненное смещение будет зафиксировано автоматически.
Для этого варианта гарантии смещение следует сохранять только после успешной обработки сообщения.
Примечание: фиксируются только смещения, превышающие текущее смещение. Например, если последнее зафиксированное смещение равно 10, и приложение выполняет offsets_store () со смещением 9, это смещение не будет зафиксировано.
В примере в предыдущем разделе Вы получили доставку «хотя бы один раз», поскольку фиксация следует за обработкой сообщения. Изменяя заказ и выполняя фиксацию синхронно перед обработкой, Вы можете получить доставку «не более одного раза», но в этом случае Вы должны тщательно обрабатывать все ошибки фиксации.
def consume_loop(consumer, topics):
try:
consumer.subscribe(topics)
while running:
msg = consumer.poll(timeout=1.0)
if msg is None: continue
if msg.error():
if msg.error().code() == KafkaError._PARTITION_EOF:
# End of partition event
sys.stderr.write('%% %s [%d] reached end at offset %d\n' %
(msg.topic(), msg.partition(), msg.offset()))
elif msg.error():
raise KafkaException(msg.error())
else:
consumer.commit(asynchronous=False)
msg_process(msg)
finally:
# Close down consumer to commit final offsets.
consumer.close()
Для простоты в этом примере перед обработкой сообщения используется Consumer.commit (). Фиксация каждого сообщения на практике приведет к большим накладным расходам. Лучшим подходом было бы собрать пакет сообщений, выполнить синхронную фиксацию, и только затем обработать сообщения (только в том случае, если фиксация прошла успешно).
Асинхронные коммиты
def consume_loop(consumer, topics):
try:
consumer.subscribe(topics)
msg_count = 0
while running:
msg = consumer.poll(timeout=1.0)
if msg is None: continue
if msg.error():
if msg.error().code() == KafkaError._PARTITION_EOF:
# End of partition event
sys.stderr.write('%% %s [%d] reached end at offset %d\n' %
(msg.topic(), msg.partition(), msg.offset()))
elif msg.error():
raise KafkaException(msg.error())
else:
msg_process(msg)
msg_count += 1
if msg_count % MIN_COMMIT_COUNT == 0:
consumer.commit(asynchronous=True)
finally:
# Close down consumer to commit final offsets.
consumer.close()
В этом примере консюмер отправляет запрос и немедленно возвращает его, используя асинхронные фиксации. Асинхронный параметр commit () изменяется на True. Значение передается явно, но асинхронные фиксации используются по умолчанию, если параметр не включен.
API возвращает обратный вызов, который вызывается при успешном или неуспешном выполнении фиксации.
from confluent_kafka import Consumer
def commit_completed(err, partitions):
if err:
print(str(err))
else:
print("Committed partition offsets: " + str(partitions))
conf = {'bootstrap.servers': "host1:9092,host2:9092",
'group.id': "foo",
'default.topic.config': {'auto.offset.reset': 'smallest'},
'on_commit': commit_completed}
consumer = Consumer(conf)
Соответствие FIPS
Confluent протестировал соответствие FIPS для клиента, использующего OpenSSL 3.0. Для использования клиента в режиме, совместимом с FIPS, Confluent рекомендует использовать OpenSSL 3.0. Более старые версии OpenSSL не были проверены (хотя они могут работать). Confluent проверил связь между клиентами и следующими конечными точками:
- Брокеры Kafka
- Регистр схемы
Брокер Kafka и Регистр схемы
- Брокер Kafka
Для общения с брокером Кафки клиент пользуется библиотекой librdkafka, которая в свою очередь использует OpenSSL. Шаги, описанные ниже, настраивают OpenSSL таким образом, чтобы можно было работать с FIPS для обеспечения связи клиента с брокером.
- Регистр схемы
Для связи с регистром схемы клиент использует стандартные библиотеки на Python, который использует родную для операционной системы библиотеку SSL/TLS. Для связи между реестром схемы и клиентом, совместимой с FIPS, не используйте следующие шаги. Вместо этого сделайте библиотеку SSL/TLS совместимой с FIPS. Если родной библиотекой SSL/TLS является OpenSSL (по умолчанию для Python), то используйте шаги в файле OpenSSL readme, чтобы сделать OpenSSL FIPS совместимым.
Коммуникация с брокером Kafka в соответствии с FIPS
Есть два способа установки клиента:
- Использовать предварительно собранные колесики
- Построить librdkafka и клиента из источника
Использование предварительно собранных колесиков
Если Вы устанавливаете этого клиента через предварительно собранные колесики с помощью confluent_kafka установки pip, то знайте, что OpenSSL 3.0 уже связан с общей библиотекой librdkafka. Чтобы разрешить этому клиенту взаимодействовать с кластером Kafka с помощью провайдера OpenSSL FIPS и алгоритмов, одобренных FIPS, необходимо включить провайдера FIPS.
ПРИМЕЧАНИЕ:
Следует включить провайдер FIPS (используя ту же самую инструкцию), если Вы устанавливаете этот клиент из исходного кода с помощью pip install confluent_kafka --no-binary: all: with prebuilt librdkafka в котором OpenSSL статически связан.
Построение librdkafka и клиента из источника
При построении librdkafka из исходного кода librdkafka динамически связывается с OpenSSL, присутствующим в системе. Если в системе установлен OpenSSL, уже работающий в режиме FIPS, то можно сразу перейти к разделу Конфигурация клиента, что позволит включить FIPS провайдера.
Если у Вас нет OpenSSL, работающего в режиме FIPS, используйте шаги, упомянутые в разделе Использование провайдера FIPS для того, чтобы сделать OpenSSL совместимым с FIPS, а затем включите провайдера fips. Как только OpenSSL начнет работать в режиме FIPS и будет включен провайдер fips, librdkafka и клиент python будут использовать одобренные FIPS алгоритмы для установлениясвязи между клиентом и кластером Kafka.
Использование провайдера FIPS
Для использования провайдера FIPS в системе должен быть доступен модуль FIPS. Подключите модуль к OpenSSL, а затем настройте OpenSSL на его использование.
Вы можете подключить провайдера FIPS в OpenSSL двумя способами:
- Поместить модуль в папку модуля по умолчанию OpenSSL.
- Указать на модуль с помощью переменной среды, OPENSSL_MODULES. Например: OPENSSL_MODULES="/path/to/fips/module/lib/folder/.
После подключения модуля провайдера FIPS к OpenSSL необходимо настроить OpenSSL. Еще раз, у вас есть два варианта:
- Изменить файл конфигурации по умолчанию, включив в него конфигурационный элемент, связанный с FIPS.
- Создать новый файл конфигурации и указать на него, используя переменную среды, OPENSSL_CONF. Например: OPENSSL_CONF="/path/to/fips/enabled/openssl/config/openssl.cnf.
ПРИМЕЧАНИЕ:
При установке клиента необходимо указать как OPENSSL_MODULES, так и OPENSSL_CONF.
Построение модуля провайдера FIPS
В этом разделе представлен общий обзор процесса построения модуля провайдера FIPS.
Построение модуля провайдера FIPS:
- Клонируйте OpenSSL из: OpenSSL Github Repo.
- Для проверки версии OpenSSL 3.0., соответствующей FIPS, используйте git checkout. Самая последняя версия может не соответствовать FIPS. На момент написания этой статьи версия 3.0.8 (помеченная как v3.0.8) является текущей протестированной версией, полностью совместимой с FIPS.
- Запустите: ./Configure enable-fips.
- Запустите: make install_fips.
В каталоге providers будут сгенерированы 2 файла, который необходимо использовать совместно с OpenSSL:
- Модуль FIPS (fips.dylib в Mac, fips.so в Linux и fips.dll в Windows)
- Конфигурация FIPS (fipsmodule.cnf)
Провайдер FIPS в OpenSSL
При установке из исходного кода можно динамически подключить встроенный выше модуль FIPS в OpenSSL, поместив модуль FIPS в папку модулей OpenSSL по умолчанию. Ищите что-то вроде:... lib/ossl-modules/.
Вы также можете указать на этот модуль с помощью переменной среды OPENSSL_MODULES.
Например: OPENSSL_MODULES="/path/to/fips/module/lib/folder/.
Запуск провайдера FIPS с помощью OpenSSL
Для включения FIPS в OpenSSL необходимо включить fipsmodule.cnf в файл, openssl.cnf. См. следующий пример openssl.cnf:
config_diagnostics = 1 openssl_conf = openssl_init .include /usr/local/ssl/fipsmodule.cnf [openssl_init] providers = provider_sect alg_section = algorithm_sect [provider_sect] fips = fips_sect [algorithm_sect] default_properties = fips=yes . . .
Файл fipsmodule.cnf содержит fips_sect, необходимые OpenSSL для включения FIPS.
Некоторые алгоритмы могут иметь другую реализацию в FIPS или в других провайдерах. При загрузке двух разных провайдеров, таких как default и fips, можно использовать любую доступную реализацию. Чтобы убедиться в том, что Вы действительно выбираете только FIPS-совместимую версию алгоритма, используйте свойство по умолчанию fips = yes в файле конфигурации.
Конфигурация клиента для запуска провайдера FIPS
OpenSSL также требует некоторых алгоритмов, не относящихся к криптографии. Эти алгоритмы не включены в провайдер FIPS, и Вам необходимо использовать базового провайдера совместно с провайдером fips. Базовый провайдер поставляется с OpenSSL по умолчанию. Необходимо включить базовый поставщик в конфигурации клиента.
Чтобы сделать клиента (клиент потребителя, производителя или администратора) совместимым с FIPS, необходимо включить в клиенте поставщика FIPS, а также базового поставщика с помощью свойства конфигурации ssl.providers.
Настройте свойство следующим образом: 'ssl.providers': 'fips, base'.




