Перейти к основному содержимому
Версия: 6.1

Интеграция Apache Kafka с приложениями Smart Beat и Logstash: совместимость и рекомендации

Статья описывает матрицу совместимости, требования, типовые конфигурации и рекомендации по эксплуатации Apache Kafka в качестве транспортного слоя для приложений Smart Beat и Logstash.

Матрица совместимости​

КомпонентПоддерживаемые версииПримечание
Apache Kafka2.8 - 4.0+Рекомендуется 3.7+ с режимом KRaft (ZooKeeper объявлен устаревшим)
Beats8.10 - 8.15output.kafka стабильно работает с Kafka 3.x/4.x
Logstash8.10 - 8.15Плагины logstash-input-kafka и logstash-output-kafka используют официальный Java-клиент Kafka

Базовые требования к инфраструктуре​

  1. Режим кластера Kafka: KRaft (рекомендуется). ZooKeeper-режим не рекомендуется для новых развертываний
  2. Сетевые требования: TCP 9092 (или другой порт) для plaintext/TLS, 9093 для SASL, если используется разделение listeners
  3. Размеры сообщений: 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)​

Используется для получения данных от сервера 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}";'
}
}
ПараметрТипЗначение по умолчаниюОписание
codecCodec"plain"Кодек, используемый для входных данных. Позволяет декодировать данные до их обработки в pipeline.
poll_timeout_msNumber100Определяет максимальное время ожидания новых данных при каждом вызове poll()
session_timeout_msNumber10000Задает максимальное время (в миллисекундах), после которого брокер считает консьюмера недоступным и инициирует ребаланс группы потребителей.
heartbeat_interval_msNumber3000Задает интервал (в миллисекундах) отправки heartbeat-сигналов от консьюмера к координатору группы потребителей. Значение не должно составлять не более трети от session_timeout_ms для своевременного подтверждения активности и предотвращения ложных ребалансов.
max_poll_recordsNumber500Задает максимальное количество сообщений, возвращаемых консьюмером за один вызов poll(). Позволяет контролировать размер обрабатываемого батча.
security_protocolString"PLAINTEXT"Задает используемый протокол безопасности. Значение может быть одним из: PLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL.
ssl_truststore_locationPath-Путь JKS Truststore для проверки сертификата брокера Kafka.
ssl_truststore_passwordPassword-Пароль от JKS Truststore, указанного в параметре ssl_truststore_location.
sasl_mechanismString"GSSAPI"Механизм SASL, используемый для установления защищенного соединения с брокерами Kafka. Определяет способ проверки учетных данных.
sasl_jaas_configString-Задает строку конфигурации JAAS, содержащую учетные данные и настройки модуля аутентификации для выбранного SASL-механизма. Является обязательным при включенной SASL-авторизации и определяет способ передачи логинов, паролей или ключей Kerberos для безопасного подключения к кластеру Kafka.

C описанием остальных параметров можно ознакомиться в документации.

При включении параметра 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
MessageSizeTooLargeExeptionmessage.max.bytes $\neq$ max.request.sizeСинхронизировать лимиты на брокере и клиенте, добавить валидацию размера в pipeline
SSL Handshake FailureНесоответствие версий TLS или отсутствие SAN в сертификатеПроверить openssl s_client, обновить сертификаты, включить ssl.verify.hostnames: false только для тестов