Конфигурирование потребителя (consumer)#

Минимальная конфигурация#

Пример конфигурационного файла consumer.properties для одного брокера (в примере в качестве хоста потребителя указан текущий хост — hostname -f):

group.id=test-consumer-group

bootstrap.servers=`hostname -f`:9101, `hostname -f`:9102,
`hostname -f`:9103

Шаги создания конфигурационного файла:

  1. Присвойте ключу group.id необходимый идентификатор группы, под которым будет запущен consumer. Можно запускать несколько consumers под одним и тем же group.id, что позволит считывать данные в несколько потоков (с нескольких серверов), обеспечивая отсутствие дублирования вычитываемых данных.

    Пример

    group.id=test-consumer-group

  2. Для ключа bootstrap.servers укажите IP:PORT сервиса Corax.

    Пример

    Если запущены 3 экземпляра Corax на том же сервере:

    bootstrap.servers=`hostname -f`:YYYY, `hostname -f`:UUUU, `hostname -f`:ZZZZ
    

    Допускается также указывать IP-адреса, полные доменные имена.

Расширенная конфигурация#

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

Запуск каждого отдельного потребителя Corax выполняется командой:

./bin/kafka-console-consumer --bootstrap-server `hostname -f`:YYYY, `hostname -f`:UUUU, `hostname -f`:ZZZZ --topic test1 --consumer.config etc/kafka/consumer.properties --from-beginning

Чтение сообщений в формате MessagePack#

Ограничение для топика

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

Для чтения сообщений в формате MessagePack необходимо запустить утилиту kafka-console-consumer.sh с указанием десериализатора сообщений (опция --value-deserializer), по умолчанию поставляется класс для работы с сообщениями от Platform V GraDeLy:

kafka-console-consumer.sh \
    --bootstrap-server <сервер> \
    --topic <топик сообщениями в MessagePack формате> \
    --consumer.config <файл с настройками> \
    --value-deserializer ru.gradely.serialization.TrailFileDeserializer

Параметры конфигурации#

Название

Описание

Тип

Значение по умолчанию

key.deserializer

Класс десериализатора для ключа, реализующего org.apache.kafka.common.serialization.Deserializerinterface

class

value.deserializer

Класс десериализатора для значения, реализующего org.apache.kafka.common.serializationDeserializerinterface

class

bootstrap.servers

Список пар «хост: порт», используемых для установления начального соединения с кластером Platform V Corax. Клиент будет использовать все серверы независимо от того, какие серверы указаны здесь для начальной загрузки. Этот список влияет только на начальные хосты, используемые для обнаружения полного набора серверов. Этот список должен быть в виде host1: port1, host2:port2,… Поскольку эти серверы используются только для первоначального подключения, чтобы обнаружить полное членство в кластере (которое может изменяться динамически), этот список не должен содержать полный набор серверов (разрешено указать более одного, на случай, если указанный сервер не работает)

list

«»

fetch.min.bytes

Минимальный объем данных, который сервер должен вернуть для запроса на выборку (fetch request). Если данных недостаточно, брокер будет ждать, пока накопится необходимое количество данных, прежде чем ответить на запрос. Значение по умолчанию 1 байт означает, что запросы на выборку отвечают, как только доступен один байт данных или время ожидания запроса на выборку данных заканчивается. Установка этого значения больше 1 заставит сервер ждать накопления больших объемов данных, что может немного повысить пропускную способность сервера за счет некоторой дополнительной задержки

int

1

group.id

Уникальная строка, идентифицирующая группу потребителей, к которой принадлежит данный потребитель. Это свойство является обязательным, если потребитель использует функции управления группами с помощью subscribe (topic) или стратегии управления смещением на основе Corax

string

null

heartbeat.interval.ms

Ожидаемое время между heartbeats (передачей сигнала «я жив») координатору потребителя при использовании средств управления группы Kafka. Передача сигнала используется для обеспечения активного сеанса потребителя и облегчения перебалансировки, когда новые потребители присоединяются или покидают группу. Значение должно быть установлено ниже session.timeout.ms, однако обычно должно устанавливается не выше 1/3 от session.timeout.ms. Его можно уменьшать с целью большего контроля предполагаемого времени для нормальных перебалансировок

int

3000

max.partition.fetch.bytes

Максимальный объем данных на партицию, который будет возвращен сервером. Записи извлекаются потребителем пакетами (batch). Если первый пакет записей в первом непустом разделе выборки больше этого предела, пакет все равно будет возвращен, чтобы гарантировать, что потребитель может добиться успеха по извлечению данных. Максимальный размер пакета записей, принятый брокером, определяется посредством message.max.bytes (настройки брокеров) или max.message.bytes (настройки топиков). Смотри fetch.max.bytes для ограничения размера запроса потребителя

int

1048576

session.timeout.ms

Тайм-аут, используемый для обнаружения сбоев потребителей при использовании средства управления группами Corax. Потребитель посылает периодические heartbeats, чтобы сообщить брокеру, что он жив (не завис). Если до истечения этого тайм-аута брокер не получит heartbeat, то брокер удалит этого потребителя из группы и инициирует перебалансировку. Значение должно находиться в допустимом диапазоне, настроенном в конфигурации брокера group.min.session.timeout.ms и group.max.session.timeout.ms

int

10000

ssl.key.password

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

password

null

ssl.keystore.location

Расположение файла хранилища ключей. Необязательный параметр для клиента и может использоваться для двусторонней проверки подлинности (authentication)

string

null

ssl.keystore.password

Пароль хранилища для файла хранилища ключей. Необязательный параметр для клиента и требуется только в случае, если объявлено ssl.keystore.location

password

null

ssl.truststore.location

Расположение файла хранилища доверенных сертификатов

string

null

ssl.truststore.password

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

password

null

auto.offset.reset

Определяет действие, выполняемое, если в Corax нет начального смещения или если текущее смещение больше не существует на сервере (например, потому что эти данные были удалены):
1) earliest — автоматический сброс смещения до самого раннего смещения;
2) latest — автоматический сброс смещения до последнего смещения;
3) none — исключение для потребителя, если для группы потребителя не найдено предыдущего смещения;
4) иное — вызов исключения потребителю

string

latest

client.dns.lookup

Управляет тем, как клиент использует DNS-запросы.
1) Если установлено значение use\all\dns\ips, то при поиске возвращаются несколько IP-адресов для имени хоста, все они будут пытаться подключиться до сбоя подключения. Применяется как к bootstrap, так и к advertised серверов.
2) Если значение resolve\canonical\bootstrap\servers\only, каждая запись будет разрешена (resolved) и расширена (expand) в список канонических имен

string

default

connections.max.idle.ms

Закрытие незанятых соединений после указанного количества миллисекунд

long

540000

default.api.timeout.ms

Задает время ожидания (в миллисекундах) для API-интерфейсов потребителя, которые могут блокировать. Эта конфигурация используется в качестве тайм-аута по умолчанию для всех операций потребителя, которые явно не принимают timeoutparameter

int

60000

enable.auto.commit

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

boolean

true

exclude.internal.topics

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

boolean

true

fetch.max.bytes

Максимальный объем данных, который сервер должен вернуть в ответ на запрос по выборке. Записи извлекаются потребителем пакетами, и если первый пакет записей в первом непустом разделе выборки больше этого значения, пакет записей все равно будет возвращен, чтобы гарантировать, что потребитель работоспособен и может извлекать данные (есть все разрешения и д.р.). Таким образом, это не абсолютный максимум. Максимальный размер пакета записей, принятый брокером, определяется посредством message.max.bytes (настройки брокеров) или max.message.bytes (настройки топиков). Обратите внимание, что потребитель выполняет несколько выборок параллельно

int

52428800

isolation.level

Управляет чтением сообщений, написанных транзакционно. Если установлено значение read\committed, то consumer.poll() будет возвращать только транзакционные сообщения, которые были зафиксированы. Если установлено значение read\uncommitted (по умолчанию), то consumer.poll() вернет все сообщения (даже транзакционные сообщения, которые были прерваны). Нетранзакционные сообщения будут возвращены безоговорочно в любом режиме. Сообщения всегда возвращаются в порядке смещения. Следовательно, в режиме read\committed consumer.poll() будет возвращать сообщения только до последнего стабильного смещения (LSO — LastStableOffset), которое меньше смещения первой открытой транзакции. В частности, любые сообщения, появляющиеся после сообщений, относящихся к текущим транзакциям, будут удерживаться до завершения соответствующей операции. В результате read\committed-потребители не смогут читать данные до high watermark, когда есть текущая транзакция (inflight). Кроме того, если установлено значение read\committed, метод seekToEnd вернет LSO

string

read\uncommitted

max.poll.interval.ms

Максимальная задержка между вызовами poll () при использовании управления группами потребителей. Это накладывает ограничение на количество времени, которое потребитель может простаивать до получения других записей. Если poll () не вызывается до истечения этого тайм-аута, то потребитель считается неработоспособным и группа будет перебалансироваться, чтобы переназначить разделы другому члену

int

300000

max.poll.records

Максимальное

int

500

partition.assignment.strategy

Имя класса стратегии назначения разделов, которую клиент будет использовать для распределения партиций между экземплярами-потребителями при использовании управления группами

list

class org.apache.kafka.clients.consumer.RangeAssignor

receive.buffer.bytes

Размер буфера приема TCP (SO\RCVBUF), используемый при чтении данных. Если значение равно -1, по умолчанию будет использоваться размер буфера, установленный в ОС

int

65536 (64 кБ)

request.timeout.ms

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

int

30000

sasl.client.callback.handler.class

Полное имя класса обработчика обратного вызова клиента SASL, реализующего интерфейс AuthenticateCallbackHandler

class

null

sasl.jaas.config

Параметры контекста входа (login) JAAS для SSL-соединений в формате, используемом файлами конфигурации JAAS. Формат файла конфигурации JAAS описан здесь. Формат значения: „loginModuleClass controlFlag (optionName=optionValue)*;“. Для брокеров конфигурация должна иметь префикс слушателя и имя механизма SASL в нижнем регистре. Например: listener.name.sasl\ssl.scram-sha-256.sasl.jaas.config=com.example.ScramLoginModule required

password

null

sasl.kerberos.service.name

Имя принципала Kerberos, под которым работает Rfarf. Это можно определить либо в конфигурации JAAS Corax, либо в конфигурации Kafka

string

null

sasl.login.callback.handler.class

Полное имя класса обработчика обратного вызова входа SASL, реализующего интерфейс обработчика AuthenticateCallbackHandler. Для брокеров конфигурация logincallbackhandler должна иметь префикс прослушивателя и имя механизма SASL в нижнем регистре. Например: listener.name.sasl\sasl.scram-sha-256.sasl.login.callback.handler.class=com.example.CustomScramLoginCallbackHandler

class

null

sasl.login.class

Полное имя класса, реализующего интерфейс входа в систему. Для брокеров конфигурация входа должна иметь префикс прослушивателя и имя механизма SASL в нижнем регистре. Например: listener.name.sasl\ssl.scram-sha-256.sasl.login.class=com.example.CustomScramLogin

class

null

sasl.mechanism

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

string

GSSAPI

security.protocol

Протокол, используемый для связи с брокерами. Допустимые значения: PLAINTEXT, SASL, SASL\PLAINTEXT, SASL\SSL

string

PLAINTEXT

send.buffer.bytes

Размер буфера отправки TCP (SO\SNDBUF), используемый при отправке данных. Если значение равно -1, по умолчанию будет использоваться параметр, установленный в ОС

int

131072

ssl.enabled.protocols

Список протоколов, разрешенных для SSL-соединений

list

TLSv1.3,TLSv1.2

ssl.keystore.type

Формат файла хранилища ключей. Параметр необязателен для клиента

string

JKS

ssl.protocol

Протокол SSL, используемый для создания SSLContext. По умолчанию — TLS, который подходит для большинства случаев. Допустимыми значениями в современных JVM являются TLSv1.2 и TLSv1.3. SSL, SSLv2 и SSLv3 могут поддерживаться в более старых JVMs, но их использование не рекомендуется из-за известных уязвимостей безопасности

string

TLS

ssl.provider

Имя поставщика безопасности, используемого для SSL-соединений. Значением по умолчанию является поставщик безопасности по умолчанию для JVM

string

null

ssl.truststore.type

Формат файла хранилища доверительных сертификатов

string

JKS

auto.commit.interval.ms

Частота в миллисекундах, с которой потребитель смещает offsets автоматически в Corax, если enable.auto.commit=true

int

5000

check.crcs

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

boolean

true

client.id

Строка id для передачи серверу при выполнении запросов. Цель этого состоит в том, чтобы иметь возможность отслеживать источник запросов за пределами только ip/port, позволяя строковое имя логического приложения включать в лог-файл запросов на стороне сервера

string

«»

fetch.max.wait.ms

Максимальное количество времени, на которое сервер заблокируется перед ответом на запрос fetch, если данных недостаточно для немедленного удовлетворения требования fetch.min.bytes

int

500

interceptor.classes

Список классов для использования в качестве перехватчиков. Реализация интерфейса org.apache.kafka.clients.ConsumerInterceptor позволяет перехватывать (и, возможно, изменять) записи, полученные потребителем. По умолчанию перехватчиков нет

list

«»

metadata.max.age.ms

Период времени в миллисекундах, после которого необходимо обновить метаданные, даже если нет каких-либо изменений лидеров партиций, чтобы проактивно обнаружить новые брокеры или партиции

long

300000

metric.reporters

Список классов для использования в качестве репортеров метрик. Реализация org.apache.kafka.common.metrics.MetricsReporterinterface позволяет подключать классы, которые будут уведомлены о создании новой метрики. JmxReporter всегда включен для регистрации статистики JMX

list

«»

metrics.num.samples

Количество выборок, сохраняемых для вычисления метрик

int

2

metrics.recording.level

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

string

INFO

metrics.sample.window.ms

Окно времени, в котором вычисляется выборка метрик

long

30000

reconnect.backoff.max.ms

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

long

1000

reconnect.backoff.ms

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

long

50

retry.backoff.ms

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

long

100

sasl.kerberos.kinit.cmd

Путь к команде Kerberos kinit

string

/usr/bin/kinit

sasl.kerberos.min.time.before.relogin

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

long

60000

sasl.kerberos.ticket.renew.jitter

Процент случайного джиттера, добавленного ко времени обновления

double

0.05

sasl.kerberos.ticket.renew.window.factor

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

double

0.8

sasl.login.refresh.buffer.seconds

Количество времени буфера до истечения срока действия учетных данных для поддержания при обновлении учетных данных в секундах. Если обновление произойдет ближе к истечению, чем количество секунд буфера, обновление будет перемещено, чтобы сохранить как можно больше времени буфера. Допустимые значения находятся в диапазоне от 0 до 3600 (1 час); если значение не указано, используется значение по умолчанию — 300 (5 минут). Это значение и значение sasl.login.refresh.min.period.seconds игнорируются, если их сумма превышает оставшееся время жизни учетных данных. В настоящее время применяется только к OAUTHBEARER

short

300

sasl.login.refresh.min.period.seconds

Требуемое минимальное время ожидания потока обновления имени входа перед обновлением учетных данных в секундах. Допустимые значения находятся в диапазоне от 0 до 900 (15 минут); если значение не указано, используется значение по умолчанию — 60 (1 минута). Это значение и значение sasl.login.refresh.buffer.seconds игнорируются, если их сумма превышает оставшееся время жизни учетных данных. В настоящее время применяется только к OAUTHBEARER

short

60

sasl.login.refresh.window.factor

Поток обновления входа (login) будет спать до тех пор, пока не будет достигнут указанный коэффициент окна относительно времени жизни учетных данных, после чего он попытается обновить учетные данные. Допустимые значения: от 0,5 (50%) до 1,0 (100%) включительно; значение по умолчанию — 0.8 (80%)

double

0.8

sasl.login.refresh.window.jitter

Максимальное количество случайного джиттера относительно времени жизни учетных данных, добавленного к времени сна потока обновления входа. Допустимые значения: от 0 до 0,25 (25%) включительно; значение по умолчанию — 0,05 (5%). В настоящее время применяется только к OAUTHBEARER

double

0.05

ssl.cipher.suites

Список шифровальных наборов. Это именованная комбинация аутентификации, шифрования, MAC и алгоритма обмена ключами, используемая для согласования параметров безопасности для сетевого подключения с использованием сетевого протокола TLS или SSL. По умолчанию поддерживаются все доступные наборы шифров

list

null

ssl.endpoint.identification.algorithm

Алгоритм идентификации конечной точки(endpoint) для проверки имени хоста сервера с помощью сертификата сервера

string

https

ssl.keymanager.algorithm

Алгоритм, используемый фабрикой менеджера ключей для SSL-соединений. Значение по умолчанию — это алгоритм key manager factory, настроенный для виртуальной машины Java

string

SunX509

ssl.secure.random.implementation

Реализация SecureRandom PRNG, используемая для операций шифрования SSL

string

null