Skip to main content
Version: 6.1

Apache Kafka Integration with SAF Beat and Logstash: Compatibility and Recommendations

This article describes compatibility, requirements, common configurations, and operational recommendations for Apache Kafka as a transport layer for SAF Beat and Logstash.

Compatibility Matrix​

ComponentSupported VersionsNotes
Apache Kafka2.8 - 4.0+3.7+ with KRaft mode is recommended; ZooKeeper is deprecated
Beats8.10 - 8.15output.kafka works reliably with Kafka 3.x/4.x
Logstash8.10 - 8.15Kafka plugins use the official Kafka Java client

Basic Infrastructure Requirements​

  1. Kafka cluster mode: KRaft is recommended. ZooKeeper mode is not recommended for new deployments
  2. Network requirements: TCP 9092 for plaintext/TLS and 9093 for SASL when listeners are separated
  3. Message size: message.max.bytes in Kafka $\geq$ max_message_bytes in Beats/Logstash. Recommended value: 10MB-20MB

Beats Configuration (Kafka Output)​

Typical Configuration​

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

Use the Kafka built-in tools for testing:

$KAFKA_HOME/bin/kafka-console-producer.sh --bootstrap-server kafka1:9092 --topic test

Beats Recommendations​

  • compression: lz4 or snappy reduces network and broker load by 30-60%
  • required_acks: 1 balances speed and reliability. Use all only when message durability is strictly required
  • configure bulk_max_size and timeout according to throughput
  • Beats does not support exactly-once. Enable idempotent mode in Kafka to minimize duplicates

Logstash Configuration (Kafka Input/Output)​

Use this plugin to receive data from Kafka. It has no required parameters.

Important

Do not use default values for production connections. Explicitly configure broker addresses, topics, consumer group, and security parameters.

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}";'
}
}

For descriptions of Kafka input parameters, see the documentation.

With decorate_events, you can extract the topic, partition, offset, and message key from Kafka metadata.

Logstash Recommendations​

  • the number of consumer_threads must not exceed the number of topic partitions, otherwise some threads remain idle
  • use decorate_events => true in the input section for debugging

Known Limitations and Workarounds​

ProblemCauseSolution
Frequent rebalancesmax.poll.interval.ms < batch processing timeIncrease the timeout, reduce max_poll_records, and optimize the pipeline
Duplicate eventsNetwork timeouts and retriesEnable an idempotent producer and use @metadata.kafka.offset for deduplication
MessageSizeTooLargeExeptionmessage.max.bytes $\neq$ max.request.sizeSynchronize broker and client limits and add pipeline size validation
SSL Handshake FailureTLS version mismatch or missing certificate SANCheck openssl s_client, update certificates, and use ssl.verify.hostnames: false only for tests