Kafka
Kafka Source
Collect logs, metrics, and traces from Apache Kafka topics using a consumer group–based ingestion model.
This source allows worker to act as a Kafka consumer, reading observability data from one or more topics and integrating Kafka-native streaming pipelines into centralized telemetry architectures.
Collection Model
- Worker connects to Kafka using bootstrap servers
- One or more topics are consumed via a consumer group
- Partitions are assigned dynamically by Kafka
- Events are decoded from raw message payloads
- Successfully processed events advance committed offsets
This model aligns with Kafka-native scalability and fault tolerance patterns.
Consumer Group Semantics
- Each worker instance joins the specified consumer group
- Kafka guarantees:
- Partition ownership by a single consumer at a time
- Automatic rebalancing on scale-up or failure
- Offsets are committed periodically based on configuration
This enables:
- Horizontal scaling
- High availability
- Exactly-once processing intent with at-least-once delivery guarantees
Offset Management and Reliability
Offset Commit Behavior
- Offsets are committed at configurable intervals
- Commit timing balances:
- Throughput
- Recovery speed
- Duplicate processing risk
On restart or rebalance:
- Consumption resumes from the last committed offset
- Uncommitted messages may be reprocessed
Shutdown and Rebalance Handling
During shutdown or partition revocation:
- Worker attempts to drain pending acknowledgements
- A configurable timeout ensures:
- Graceful shutdown
- Consumer group stability
- Prevention of forced rebalance eviction
Topic Selection Strategy
- Topics can be defined explicitly
- Regular expressions are supported
- Enables:
- Dynamic topic discovery
- Prefix-based ingestion
- Multi-tenant Kafka usage
This is particularly useful in environments with high topic churn.
Event Decoding
Kafka messages are delivered as raw bytes and must be decoded.
Supported Codecs
- JSON
- Syslog (RFC 3164 / RFC 5424)
- GELF
- Avro (including schema ID stripping)
- Protobuf
- OTLP (logs, metrics, traces)
- InfluxDB Line Protocol
- Worker native formats
- Raw bytes
- VRL-based custom decoding
Some codecs can automatically infer the event signal type.
Multi-Signal Support
The Kafka source supports:
- Logs
- Metrics
- Traces
For OTLP and native formats:
- Signal types can be auto-detected
- Parsing order can be restricted for performance optimization
Framing and Message Boundaries
Framing defines how events are extracted from Kafka message payloads.
Supported strategies include:
- Byte passthrough
- Newline-delimited
- Character-delimited
- Length-delimited
- Octet counting
- Varint length-delimited
- Chunked GELF
Frame size limits prevent malformed messages from exhausting memory.
Metadata Enrichment
Kafka-specific metadata can be attached to each event, including:
- Topic name
- Partition number
- Offset
- Message key
- Message headers
These fields enable:
- Debugging
- Replay logic
- Downstream routing and filtering
- Auditing and lineage tracking
Kafka Client Configuration
librdkafka Integration
Worker uses librdkafka internally, enabling:
- Fine-grained client tuning
- Broker-specific optimizations
- Advanced protocol features
Low-level Kafka parameters can be passed directly when needed.
Security and Authentication
TLS Support
- Encrypted connections to Kafka brokers
- Custom CA bundles
- Certificate verification and hostname validation
- Required for production-grade clusters
SASL Authentication
Supported mechanisms:
- PLAIN
- SCRAM-based authentication
For advanced mechanisms (e.g. Kerberos):
- Native librdkafka options are used
This allows compatibility with enterprise Kafka deployments.
Metrics and Observability
Consumer Lag Metrics
Optional metrics expose:
- Per-topic lag
- Per-partition lag
These metrics are critical for:
- Backlog detection
- Capacity planning
- Pipeline health monitoring
Reliability Characteristics
- At-least-once delivery semantics
- Stream-based ingestion
- Backpressure-aware processing
- Offset-based recovery
- Stateless processing model
These properties make the Kafka source suitable for high-volume, continuous telemetry ingestion.
Common Use Cases
- Centralized log aggregation from Kafka
- Metrics and trace ingestion via OTLP over Kafka
- Streaming observability pipelines
- Multi-tenant Kafka telemetry platforms
- Decoupled producer–consumer architectures