Интеграция Apache Kafka с приложениями Smart Beat и Logstash: совместимость и рекомендации
Статья описывает матрицу совместимости, требования, типовые конфигурации и рекомендации по эксплуатации Apache Kafka в качестве транспортного слоя для приложений Smart Beat и Logstash.
Матрица совместимости
| Компонент | Поддерживаемые версии | Примечание |
|---|---|---|
Apache Kafka | 2.8 - 4.0+ | Рекомендуется 3.7+ с режимом KRaft (ZooKeeper объявлен устаревшим) |
Beats | 8.10 - 8.15 | output.kafka стабильно работает с Kafka 3.x/4.x |
Logstash | 8.10 - 8.15 | Плагины logstash-input-kafka и logstash-output-kafka используют официальный Java-клиент Kafka |
Базовые требования к инфраструктуре
- Режим кластера Kafka: KRaft (рекомендуется). ZooKeeper-режим не рекомендуется для новых развертываний
- Сетевые требования: TCP
9092(или другой порт) для plaintext/TLS,9093для SASL, если используется разделение listeners - Размеры сообщений:
message.max.bytesв Kafka $\geq$max_message_bytesв Beats/Logstash. Рекомендуемое значение:10MB-20MB
Настройка Beats (Kafka Output)
Типовая конфигурация
output.kafka:
hosts: ["kafka1:9092", "kafka2:9092", "kafka3:9092"]
topic: "smartbeat-%{[agent.name]:default}"
partition.round_robin:
reachable_only: false
required_acks: 1
compression: lz4
max_message_bytes: 1000000
Для тестов можно использовать встроенные в Kafka средства. Пример использования для добавления сообщений (запускать из папки, где установлен Apache Kafka):
$KAFKA_HOME/bin/kafka-console-producer.sh --bootstrap-server kafka1:9092 --topic test
Рекомендации для Beats
compression: lz4илиsnappyснижает нагрузку на сеть и брокеры на 30-60%required_acks: 1оптимален для баланса между скоростью и надежностью. Используйтеallтолько при строгом требовании к сохранности сообщенийbulk_max_sizeиtimeoutнастраивайте под пропускную способность- Beats не поддерживает
exactly-once. Для минимизации дубликатов включитеidempotent-режим на стороне Kafka (включен по умолчанию с версии 3.0+)
Настройка Logstash (Kafka Input/Output)
- input
- output
Используется для получения данных от сервера Kafka. Для данного плагина отсутствуют обязательные параметры.
Мы не рекомендуем использовать параметры по умолчанию. Для рабочего подключения явно задайте адреса брокеров, топики, группу потребителей и параметры безопасности. Значения по умолчанию рассчитаны на локальное подключение и могут не соответствовать вашей конфигурации.
input {
kafka {
bootstrap_servers => "kafka1:9092,kafka2:9092,kafka3:9092"
topics => ["beats-ingest", "app-logs"]
group_id => "logstash-consumer-group"
auto_offset_reset => "latest"
codec => "json"
decorate_events => true
consumer_threads => 4
poll_timeout_ms => 5000
session_timeout_ms => 30000
heartbeat_interval_ms => 10000
max_poll_records => 500
security_protocol => "SASL_SSL"
ssl_truststore_location => "/etc/logstash/kafka-truststore.jks"
ssl_truststore_password => "${TRUSTSTORE_PASS}"
sasl_mechanism => "SCRAM-SHA-512"
sasl_jaas_config => 'org.apache.kafka.common.security.scram.ScramLoginModule required username="${KAFKA_USER}" password="${KAFKA_PASSWORD}";'
}
}
| Параметр | Тип | Значение по умолчанию | Описание |
|---|---|---|---|
codec | Codec | "plain" | Кодек, используемый для входных данных. Позволяет декодировать данные до их обработки в pipeline. |
poll_timeout_ms | Number | 100 | Определяет максимальное время ожидания новых данных при каждом вызове poll() |
session_timeout_ms | Number | 10000 | Задает максимальное время (в миллисекундах), после которого брокер считает консьюмера недоступным и инициирует ребаланс группы потребителей. |
heartbeat_interval_ms | Number | 3000 | Задает интервал (в миллисекундах) отправки heartbeat-сигналов от консьюмера к координатору группы потребителей. Значение не должно составлять не более трети от session_timeout_ms для своевременного подтверждения активности и предотвращения ложных ребалансов. |
max_poll_records | Number | 500 | Задает максимальное количество сообщений, возвращаемых консьюмером за один вызов poll(). Позволяет контролировать размер обрабатываемого батча. |
security_protocol | String | "PLAINTEXT" | Задает используемый протокол безопасности. Значение может быть одним из: PLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL. |
ssl_truststore_location | Path | - | Путь JKS Truststore для проверки сертификата брокера Kafka. |
ssl_truststore_password | Password | - | Пароль от JKS Truststore, указанного в параметре ssl_truststore_location. |
sasl_mechanism | String | "GSSAPI" | Механизм SASL, используемый для установления защищенного соединения с брокерами Kafka. Определяет способ проверки учетных данных. |
sasl_jaas_config | String | - | Задает строку конфигурации JAAS, содержащую учетные данные и настройки модуля аутентификации для выбранного SASL-механизма. Является обязательным при включенной SASL-авторизации и определяет способ передачи логинов, паролей или ключей Kerberos для безопасного подключения к кластеру Kafka. |
C описанием остальных параметров можно ознакомиться в документации.
Данный плагин используется для отправки данных брокеру Kafka. Обязательный параметр: topic_id.
output {
kafka {
bootstrap_servers => "kafka1:9092,kafka2:9092,kafka3:9092"
topic_id => "processed-%{[fields][pipeline]}"
codec => "json"
compression_type => "lz4"
acks => "1"
batch_size => 16384
linger_ms => 10
max_request_size => 10485760
retries => 3
retry_backoff_ms => 1000
security_protocol => "SASL_SSL"
sasl_mechanism => "SCRAM-SHA-512"
sasl_jaas_config => 'org.apache.kafka.common.security.scram.ScramLoginModule required username="${KAFKA_USER}" password="${KAFKA_PASSWORD}";'
}
}
| Параметр | Тип | Значение по умолчанию | Описание |
|---|---|---|---|
bootstrap_servers | String | "localhost:9092" | Задает начальный список брокеров Kafka в формате "host:port". |
topic_id | String | - | Задает имя Kafka-топика, в который плагин будет публиковать исходящие события. Поддерживает динамическое формирование названия через ссылки на поля событий: logs-%{[service]}. |
codec | Codec | "plain" | Определяет способ сериализации событий Logstash перед публикацией в Kafka. |
compression_type | String | "none" | Задает алгоритм сжатия сообщений перед отправкой в Kafka. Может быть одним из значений: none, gzip, snappy, lz4, zstd. |
acks | String | "1" | Задает минимальное количество подтверждений от брокеров Kafka, необходимое продюсеру для успешной отправки сообщения. Может быть одним из значений: 0, 1, all. |
batch_size | Number | 16384 | Задает максимальный размер (в байтах) буфера, в котором продюсер накапливает сообщения перед отправкой в одну партицию Kafka. |
linger_ms | Number | 0 | Задает максимальное время ожидания перед отправкой накопленных сообщений, даже если целевой размер батча (batch_size) еще не достигнут. |
max_request_size | Number | 1048576 | Задает максимально допустимый размер (в байтах) одного сетевого запроса, отправляемого продюсером брокеру Kafka. |
retries | Number | - | Задает максимальное количество попыток повторной отправки сообщений брокеру Kafka при возникновении транзитных ошибок. |
retry_backoff_ms | Number | 100 | Задает базовую задержку между повторными попытками отправки сообщения брокеру при возникновении временных ошибок. |
security_protocol | String | PLAINTEXT | Задает базовый протокол безопасности для коммуникации с брокерами Kafka, определяя использование шифрования канала связи (SSL/TLS) и механизмов аутентификации (SASL). Может быть одним из значений: PLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL. |
sasl_mechanism | String | "GSSAPI" | Задает алгоритм SASL-аутентификации, который Logstash использует для подтверждения своей подлинности брокерам Kafka. |
sasl_jaas_config | String | - | Задает строку конфигурации JAAS, содержащую учетные данные для SASL-аутентификации Logstash при подключении к брокерам Kafka. Является обязательным при использовании SASL и обеспечивает безопасную передачу секретов, необходимых для успешной публикации сообщений в защищенные топики. |
При включении параметра decorate_events можно получить дополнительные данные сообщения из Kafka. В примере ниже извлекаются имя топика (topic), раздел (partition), смещение (offset) и ключ сообщения:
filter {
mutate{
add_field => { "[topic_name]" => "%{[@metadata][kafka][topic]}"}
add_field => { "[topic_partition]" => "%{[@metadata][kafka][partition]}"}
add_field => { "[topic_offset]" => "%{[@metadata][kafka][offset]}"}
add_field => { "[topic_key]" => "%{[@metadata][kafka][key]}"}
}
}
В фильтр можно добавить следующие параметры:
add_field => { "[topic_consumer_group]" => "%{[@metadata][kafka][consumer_group]}"}
add_field => { "[topic_timestamp]" => "%{[@metadata][kafka][timestamp]}"}
Это позволит увидеть группу потребителей, в которой выполнялось чтение сообщения, и время чтения.
Рекомендации для Logstash
- количество
consumer_threadsне должно превышать количество разделов топика, иначе часть потоков будет простаивать - Используйте
decorate_events => trueв секцииinputдля отладки (добавляет@metadata.kafka.*)
Известные ограничения и обходные пути
| Проблема | Причина | Решение |
|---|---|---|
| Частые ребалансы | max.poll.interval.ms < время обработки батча | Увеличить таймаут, уменьшить max_poll_records, оптимизировать пайплайн |
| Дубликаты событий | Сетевые таймауты + retry | Включить idempotent producer, использовать @metadata.kafka.offset для дедупликации в Logstash |
MessageSizeTooLargeExeption | message.max.bytes $\neq$ max.request.size | Синхронизировать лимиты на брокере и клиенте, добавить валидацию размера в pipeline |
| SSL Handshake Failure | Несоответствие версий TLS или отсутствие SAN в сертификате | Проверить openssl s_client, обновить сертификаты, включить ssl.verify.hostnames: false только для тестов |