Apache Kafka Integration with Smart Beat and Logstash: Compatibility and Recommendations
This article describes compatibility, requirements, common configurations, and operational recommendations for Apache Kafka as a transport layer for Smart Beat and Logstash.
Compatibility Matrix
| Component | Supported Versions | Notes |
|---|---|---|
Apache Kafka | 2.8 - 4.0+ | 3.7+ with KRaft mode is recommended; ZooKeeper is deprecated |
Beats | 8.10 - 8.15 | output.kafka works reliably with Kafka 3.x/4.x |
Logstash | 8.10 - 8.15 | Kafka plugins use the official Kafka Java client |
Basic Infrastructure Requirements
- Kafka cluster mode: KRaft is recommended. ZooKeeper mode is not recommended for new deployments
- Network requirements: TCP
9092for plaintext/TLS and9093for SASL when listeners are separated - Message size:
message.max.bytesin Kafka $\geq$max_message_bytesin 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: lz4orsnappyreduces network and broker load by 30-60%required_acks: 1balances speed and reliability. Useallonly when message durability is strictly required- configure
bulk_max_sizeandtimeoutaccording to throughput - Beats does not support
exactly-once. Enable idempotent mode in Kafka to minimize duplicates
Logstash Configuration (Kafka Input/Output)
- 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.
Use this plugin to send data to Kafka. Required parameter: 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}";'
}
}
With decorate_events, you can extract the topic, partition, offset, and message key from Kafka metadata.
Logstash Recommendations
- the number of
consumer_threadsmust not exceed the number of topic partitions, otherwise some threads remain idle - use
decorate_events => truein theinputsection for debugging
Known Limitations and Workarounds
| Problem | Cause | Solution |
|---|---|---|
| Frequent rebalances | max.poll.interval.ms < batch processing time | Increase the timeout, reduce max_poll_records, and optimize the pipeline |
| Duplicate events | Network timeouts and retries | Enable an idempotent producer and use @metadata.kafka.offset for deduplication |
MessageSizeTooLargeExeption | message.max.bytes $\neq$ max.request.size | Synchronize broker and client limits and add pipeline size validation |
| SSL Handshake Failure | TLS version mismatch or missing certificate SAN | Check openssl s_client, update certificates, and use ssl.verify.hostnames: false only for tests |