Конфигурирование потребителя (consumer)#
Минимальная конфигурация#
Пример конфигурационного файла consumer.properties для одного брокера (в примере в качестве хоста потребителя указан текущий хост — hostname -f):
group.id=test-consumer-group
bootstrap.servers=`hostname -f`:9101, `hostname -f`:9102,
`hostname -f`:9103
Шаги создания конфигурационного файла:
Присвойте ключу
group.idнеобходимый идентификатор группы, под которым будет запущен consumer. Можно запускать несколько consumers под одним и тем жеgroup.id, что позволит считывать данные в несколько потоков (с нескольких серверов), обеспечивая отсутствие дублирования вычитываемых данных.Пример
group.id=test-consumer-groupДля ключа
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
Название |
Описание |
Тип |
Значение по умолчанию |
|---|---|---|---|
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). Если первый пакет записей в первом непустом разделе выборки больше этого предела, пакет все равно будет возвращен, чтобы гарантировать, что потребитель может добиться успеха по извлечению данных. Максимальный размер пакета записей, принятый брокером, определяется посредством |
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 нет начального смещения или если текущее смещение больше не существует на сервере (например, потому что эти данные были удалены): |
string |
latest |
client.dns.lookup |
Управляет тем, как клиент использует DNS-запросы. |
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 |
Управляет чтением сообщений, написанных транзакционно. Если установлено значение |
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 |