Introduction
In today’s cloud‑native world, logs are the nervous system of every application stack. They tell us when a request succeeded, why a service failed, and how resources are being consumed. Yet as systems grow—spanning containers, serverless functions, edge devices, and even autonomous AI agents—the challenge shifts from “how do I collect logs?” to “how do I route them efficiently to every destination that needs them?”
Enter Fluentd, the open‑source data collector that has become the de‑facto standard for log aggregation and routing. Designed with a pluggable architecture, Fluentd can ingest data from dozens of sources, transform it on the fly, and fan‑out to multiple back‑ends such as Elasticsearch, Amazon S3, Kafka, or a custom AI‑driven analytics platform. For a platform like Apiary, which monitors thousands of beehives via IoT sensors and coordinates self‑governing AI agents that adjust hive conditions, Fluentd’s ability to reliably ship the same event to a monitoring dashboard, a long‑term archive, and a machine‑learning model in parallel is not a luxury—it’s a necessity.
This article is a deep dive into Fluentd as a routing powerhouse. We’ll explore its core concepts, walk through real configuration patterns for multi‑destination logging, examine performance and security considerations, and showcase concrete examples from both commercial and conservation‑focused workloads. By the end, you’ll have a production‑ready blueprint for turning raw log streams into actionable insight across any number of endpoints.
1. What Is Fluentd? Architecture and Core Components
Fluentd was created by Treasure Data in 2011 and donated to the Cloud Native Computing Foundation (CNCF) in 2015. It is written in Ruby (with a C‑based core for performance) and follows the “unified logging layer” philosophy: a single agent that can speak the language of every log source and destination.
1.1 Data Flow Model
At its heart, Fluentd processes events through a pipeline composed of four stages:
| Stage | Purpose | Typical Plugins |
|---|---|---|
| Input | Pull or receive raw logs from a source (e.g., files, syslog, HTTP) | in_tail, in_forward, in_http |
| Filter | Enrich, parse, or drop records (JSON parsing, field renaming, geo‑IP lookup) | filter_parser, filter_record_transformer, filter_grok |
| Buffer | Temporarily store events to guarantee delivery and smooth bursts | Memory, file, or chunked buffering |
| Output | Push processed events to one or more destinations | out_elasticsearch, out_s3, out_kafka, out_stdout |
Each plugin is a Ruby class that adheres to a simple interface (configure, start, shutdown, emit). The modularity means you can swap an out_s3 for an out_bigquery without touching the rest of the pipeline.
1.2 The “Label” Mechanism
Fluentd introduces labels (also called routes) to partition the configuration into independent processing graphs. A label is a named block that can contain its own inputs, filters, and outputs. When an event matches a <match> directive, it is directed to the corresponding label’s pipeline. This is the foundation for multi‑destination routing: the same event can be matched by multiple <match> statements, each sending a copy to a different label.
<label @ROUTER>
<match **>
@type copy
<store>
@type elasticsearch
host es.example.com
</store>
<store>
@type s3
bucket apiary-logs
</store>
</match>
</label>
In the example above, every record (**) is duplicated and shipped to both Elasticsearch and Amazon S3. The copy plugin handles the fan‑out without requiring the user to duplicate the entire pipeline.
1.3 Why Fluentd Over Alternatives?
| Feature | Fluentd | Fluent Bit | Logstash |
|---|---|---|---|
| Plugin count (as of 2024) | > 800 | ~ 70 | ~ 200 |
| Language | Ruby (extensible) | C (lightweight) | Java (JVM) |
| Resource footprint | 30–80 MiB per instance | 5–15 MiB | 150–300 MiB |
| Multi‑destination fan‑out | Native copy & label | Limited (requires multiple instances) | Built‑in pipeline |
| Community support | CNCF, strong enterprise adopters | CNCF, edge‑focused | Elastic ecosystem |
For workloads that need rich transformation and simultaneous delivery to several back‑ends—think sending hive sensor telemetry to a time‑series DB, a data lake, and a real‑time anomaly detector—Fluentd’s plugin ecosystem and routing semantics make it the most versatile choice.
2. Core Concepts for Multi‑Destination Routing
Understanding how Fluentd decides where to send a log entry is essential before you write any configuration. Below are the key concepts that enable precise, scalable routing.
2.1 Tagging and Wildcards
Every event in Fluentd carries a tag (a dot‑separated string). Tags are assigned by the input plugin, often derived from the file path or the source hostname. For example, an input reading /var/log/nginx/access.log might emit the tag nginx.access.
Tags can be matched using glob‑style wildcards:
nginx.*matchesnginx.accessandnginx.error**matches any tag (the catch‑all)
Using tags, you can direct only the logs you care about to a specific destination.
2.2 The <match> Directive
A <match> block defines where events with a given tag pattern should go. Inside a <match>, you specify an output plugin (or a copy plugin for multiple stores). Example:
<match nginx.access>
@type elasticsearch
host es-prod.example.com
logstash_format true
</match>
Multiple <match> blocks can target the same tag pattern, enabling parallel routing:
<match nginx.access>
@type copy
<store>
@type elasticsearch
host es-prod.example.com
</store>
<store>
@type s3
bucket nginx-logs
path logs/%Y/%m/%d/
</store>
</match>
2.3 The copy Plugin
The copy plugin is the workhorse for fan‑out. It receives an event once, then forwards a deep copy to each <store> child. This guarantees that downstream plugins cannot interfere with each other’s state.
Key parameters:
| Parameter | Description |
|---|---|
store_retry_limit | Max retries per store before discarding (default 17) |
store_flush_interval | How often to flush buffers (seconds) |
store_log_level | Log level for the copy operation |
2.4 Conditional Routing with <store>
Within a copy, you can add conditional filters to decide per‑record whether a store should receive it. The @type relabel plugin can re‑tag events and re‑enter the routing engine:
<store>
@type relabel
@label @ANOMALY
<condition>
key severity
pattern ^(error|critical)$
</condition>
</store>
If the condition matches, the event is sent to the @ANOMALY label, where a separate pipeline can forward it to a Slack webhook or an AI incident‑response service.
2.5 Buffer Types and Guarantees
Fluentd’s buffering strategy determines delivery semantics:
- Memory buffer: fastest, but events are lost on crash.
- File buffer: persists on disk; supports at‑least‑once delivery.
- Chunked buffer: groups events into chunks (default 8 MiB) for efficient bulk writes.
When routing to multiple destinations, you often want independent buffering per store to avoid a slow destination throttling the entire pipeline. The copy plugin automatically creates a separate buffer for each <store> unless you explicitly share one with @type file at the copy level.
3. Setting Up Fluentd for Multi‑Destination Logging
Below is a step‑by‑step guide to building a production‑grade Fluentd configuration that routes logs to three common back‑ends:
- Elasticsearch – real‑time search & Kibana dashboards.
- Amazon S3 – immutable, cost‑effective long‑term archive.
- Kafka – streaming platform for downstream AI agents.
3.1 Installing Fluentd
On Ubuntu 22.04:
curl -L https://toolbelt.treasuredata.com/sh/install-ubuntu-focal-td-agent4.sh | sh
sudo systemctl enable td-agent
sudo systemctl start td-agent
The package installs td-agent, the official Fluentd binary, along with a default config at /etc/td-agent/td-agent.conf.
3.2 Defining Input Sources
For a typical Apiary deployment, logs come from:
- File logs generated by the hive‑monitoring daemon (
/var/log/apiary/hive/*.log). - HTTP JSON payloads sent by edge AI agents (
/api/v1/telemetry).
# File tailing
<source>
@type tail
path /var/log/apiary/hive/*.log
pos_file /var/log/td-agent/hive.pos
tag hive.log
<parse>
@type json
</parse>
</source>
# HTTP input for AI agents
<source>
@type http
port 9880
bind 0.0.0.0
body_size_limit 32m
keepalive_timeout 10
tag ai.telemetry
</source>
Both sources emit JSON, making downstream parsing trivial.
3.3 Central Routing Block
We’ll use a label named @ROUTER that contains a copy store for each destination. Note the separate buffers to isolate performance.
<label @ROUTER>
<match hive.log>
@type copy
# Elasticsearch store
<store>
@type elasticsearch
host es-prod.example.com
port 9200
logstash_format true
flush_interval 5s
buffer_type file
buffer_path /var/log/td-agent/buffer/es
</store>
# S3 archive store
<store>
@type s3
aws_key_id YOUR_AWS_KEY
aws_sec_key YOUR_AWS_SECRET
s3_bucket apiary-logs
s3_region us-east-1
path logs/hive/%Y/%m/%d/
store_as gzip
buffer_type file
buffer_path /var/log/td-agent/buffer/s3
</store>
# Kafka stream store
<store>
@type kafka2
brokers kafka01.example.com:9092,kafka02.example.com:9092
topic hive_telemetry
use_event_time true
buffer_type file
buffer_path /var/log/td-agent/buffer/kafka
</store>
</match>
# Route AI telemetry to a different set of destinations
<match ai.telemetry>
@type copy
<store>
@type stdout
</store>
<store>
@type elasticsearch
host es-ml.example.com
logstash_format true
index_name ai-telemetry
</store>
<store>
@type file
path /var/log/apiary/ai/telemetry.log
</store>
</match>
</label>
3.4 Verifying the Pipeline
After reloading Fluentd (sudo systemctl reload td-agent), you can verify each destination:
- Elasticsearch:
curl -XGET 'http://es-prod.example.com:9200/_cat/indices?v'should showhive.log-2024.09.27indices. - S3: Use AWS CLI
aws s3 ls s3://apiary-logs/logs/hive/2024/09/27/to confirm gzipped files appear. - Kafka:
kafka-console-consumer.sh --bootstrap-server kafka01.example.com:9092 --topic hive_telemetry --from-beginningshould stream JSON records.
3.5 Scaling with Multiple Workers
Fluentd supports worker threads (-w N) to parallelize processing. For a high‑throughput hive monitoring fleet that generates ~10 k events per second, a configuration with -w 4 and --log-rotate-age 7 typically sustains the load while keeping latency under 200 ms per event.
sudo td-agent -c /etc/td-agent/td-agent.conf -w 4 -v
4. Advanced Routing Strategies
Beyond simple copy‑and‑paste, Fluentd offers sophisticated mechanisms for conditional routing, dynamic tag rewriting, and load‑balanced distribution. These patterns are especially useful when you need to route different subsets of logs to specialized AI agents or balance traffic across multiple data lakes.
4.1 Conditional Copy with @type rewrite_tag_filter
Suppose you want to send only critical hive alerts (severity critical) to a real‑time incident bot, while all other logs go to the standard pipeline.
<filter hive.log>
@type rewrite_tag_filter
rewriterule1 severity ^critical$ critical_alert
</filter>
<label @CRITICAL>
<match critical_alert>
@type copy
<store>
@type http
endpoint https://incident-bot.example.com/webhook
serializer json
</store>
<store>
@type elasticsearch
host es-incident.example.com
index_name hive-critical
</store>
</match>
</label>
The rewrite_tag_filter creates a new tag critical_alert for matching records, which are then processed by the @CRITICAL label.
4.2 Load Balancing with @type copy + store_retry_limit
When sending telemetry to a cluster of Kafka brokers, you may wish to spread partitions while also providing a fallback if a broker becomes unavailable. The copy plugin’s per‑store retry logic can be tuned:
<store>
@type kafka2
brokers kafka01.example.com:9092
topic hive_telemetry
buffer_type file
retry_wait 5s
store_retry_limit 10 # after 10 attempts, drop to next store
</store>
<store>
@type kafka2
brokers kafka02.example.com:9092
topic hive_telemetry
buffer_type file
retry_wait 5s
store_retry_limit 10
</store>
If kafka01 fails, the event is retried 10 times, then automatically handed off to kafka02.
4.3 Dynamic Destination Selection with ${tag} and ${record}
Fluentd allows runtime interpolation in output parameters. This is handy for routing logs to per‑tenant indices or buckets without hard‑coding each one.
<match hive.*>
@type copy
<store>
@type elasticsearch
host es-prod.example.com
index_name ${tag}_%Y.%m.%d # e.g., hive.log_2024.09.27
</store>
<store>
@type s3
s3_bucket apiary-${record["hive_id"]}-logs
path ${record["hive_id"]}/%Y/%m/%d/
</store>
</match>
If a hive sensor includes "hive_id": "HB-42" in its JSON payload, the S3 bucket becomes apiary-HB-42-logs, enabling per‑hive isolation for compliance or data‑ownership policies.
4.4 Multi‑Label Fan‑Out with @type label_router
For extremely complex topologies—such as sending logs to both a real‑time monitoring dashboard and a batch analytics pipeline—the label_router plugin can split a stream into multiple labels without copying the data.
<label @MAIN>
<match **>
@type label_router
<route>
@label @REALTIME
</route>
<route>
@label @BATCH
</route>
</match>
</label>
<label @REALTIME>
<match **>
@type elasticsearch
host es-realtime.example.com
logstash_format true
</match>
</label>
<label @BATCH>
<match **>
@type s3
s3_bucket apiary-batch-logs
path batch/%Y/%m/%d/
store_as gzip
</match>
</label>
Because the router does not duplicate the event, memory usage stays low, while each label can apply its own buffering strategy (e.g., a tiny memory buffer for real‑time, a larger file buffer for batch).
5. Performance, Reliability, and Scaling
Routing logs to multiple destinations can become a bottleneck if not engineered correctly. Below we examine the metrics that matter, the buffering knobs that keep data safe, and the scaling patterns proven in production.
5.1 Throughput Benchmarks
A benchmark conducted by Treasure Data in Q2 2023 measured single‑node Fluentd processing ~150 k events per second with a 4‑core CPU and 16 GiB RAM, using the following setup:
| Scenario | Plugins | Avg Latency | CPU Utilization |
|---|---|---|---|
Single copy → Elasticsearch + S3 | in_tail, filter_parser, copy (2 stores) | 120 ms | 68 % |
copy + label_router → 3 stores (Elasticsearch, Kafka, S3) | in_http, copy (3 stores) | 170 ms | 82 % |
| Multi‑worker (4 workers) | Same as above | 85 ms | 95 % (across cores) |
These numbers demonstrate that adding more workers (via -w) scales linearly until CPU saturation. For workloads exceeding 200 k events/sec, horizontal scaling—running multiple Fluentd instances behind a load balancer—is recommended.
5.2 Buffering Strategies for At‑Least‑Once Delivery
When routing to durable stores (S3, Kafka, Elasticsearch), you typically want at‑least‑once guarantees. Fluentd achieves this by:
- Chunking events into files (default 8 MiB).
- Writing the chunk to the destination.
- Acknowledging success before deleting the chunk.
If the destination returns an error, the chunk is re‑queued according to retry_wait and retry_limit.
Best practice: set flush_interval to a low value (e.g., 5s) for near‑real‑time pipelines, and flush_thread_count to 2 to allow parallel flushes.
<store>
@type s3
buffer_type file
buffer_path /var/log/td-agent/buffer/s3
flush_interval 5s
retry_wait 10s
max_retry_wait 300s
</store>
5.3 Back‑Pressure Handling
If one destination slows down (e.g., S3 experiences throttling), Fluentd’s per‑store buffers prevent the slowdown from propagating upstream. However, if all buffers fill, Fluentd will drop new events according to the overflow_action setting (block, drop_oldest_chunk, drop_newest_chunk). For critical telemetry, use overflow_action block to pause the input source (e.g., tail will stop reading the file) rather than losing data.
5.4 Horizontal Scaling with Fluent Bit Edge Agents
For edge devices such as beehive sensors, running a full‑featured Fluentd instance may be overkill. Fluent Bit—the lightweight C implementation—can act as a forwarder, sending logs to a central Fluentd aggregator. This pattern reduces per‑device CPU usage to <2 % and bandwidth by batching locally.
# On the edge device
[INPUT]
Name tail
Path /var/log/hive/telemetry.log
Tag hive.edge
[OUTPUT]
Name forward
Match *
Host fluentd-aggregator.example.com
Port 24224
The central Fluentd then performs the heavy routing, while Fluent Bit handles collection and initial buffering.
6. Real‑World Use Cases
6.1 E‑Commerce Transaction Logging
A global retailer processes **5