Apache Flume was built to answer one question for Hadoop clusters: how do you move a continuous stream of log lines from hundreds of servers into HDFS without losing them when a disk fills, a collector dies or the NameNode is slow? Its answer is small and still worth understanding: every hop is a transaction against a buffer, and data leaves a buffer only after the next hop has accepted it. That one rule gives you at-least-once delivery and tells you exactly where data can be lost when you break it.
This article explains Flume from that rule outward: events and agents, the transaction contract, the choice between memory, file and Kafka channels, a complete two-tier topology from application hosts to HDFS with failover, the HDFS sink settings that decide whether you produce healthy files or millions of tiny ones, a custom sink, monitoring and failure modes. A note on status: the latest release listed on the project site is 1.11.0, from October 2022, and many teams now choose Kafka-based pipelines for new work. Flume still runs in many Hadoop estates, so operating it well and knowing when to migrate both matter.
Events, agents and the three component types
A Flume event is a byte array body plus a map of string headers. Flume never parses the body unless you ask an interceptor or serializer to; headers carry routing and partitioning information such as a timestamp, a host name or a service name. An agent is a JVM process that hosts three kinds of components wired together in a properties file. A source receives or reads data and turns it into events. A channel is a buffer that holds events until something takes them. A sink takes events from one channel and writes them somewhere: HDFS, Kafka, another agent over Avro, or a log.
A source can write to several channels, chosen by a channel selector: replicating (every channel gets a copy, the default) or multiplexing (a header value picks the channel, with selector.header, selector.mapping.<value> and selector.default). A sink reads from exactly one channel. Sinks can be grouped under a sink processor for failover or load balancing. Interceptors sit between source and channel and can add headers (timestamp, host, static values) or drop events by regular expression.
The transaction contract and what it guarantees
Each channel exposes transactions. A source opens one, puts a batch of events, and commits; if the put fails it rolls back and the client sees an error. A sink opens one, takes a batch, writes it downstream, and commits only after the downstream write succeeded; on any failure it rolls back and the events become available again. Chain two agents with an Avro sink and an Avro source and the guarantee extends across the network: the upstream sink commits only after the downstream source has committed into its own channel. This is the loop every sink runs, written against the real Flume API:
public class AuditSink extends AbstractSink implements Configurable {
private int batchSize;
public void configure(Context ctx) { batchSize = ctx.getInteger("batchSize", 500); }
public Status process() throws EventDeliveryException {
Channel ch = getChannel();
Transaction tx = ch.getTransaction();
tx.begin();
try {
List<Event> batch = new ArrayList<>(batchSize);
for (int i = 0; i < batchSize; i++) {
Event e = ch.take();
if (e == null) break; // channel empty
batch.add(e);
}
if (batch.isEmpty()) { tx.commit(); return Status.BACKOFF; }
writeDownstream(batch); // must throw on failure
tx.commit(); // only now do events leave the channel
return Status.READY;
} catch (Throwable t) {
tx.rollback(); // events become visible again
throw new EventDeliveryException("write failed; batch will be retried", t);
} finally {
tx.close();
}
}
}The consequence is at-least-once delivery. If the downstream write succeeds but the commit is lost, because the agent crashed in between, the batch is delivered again. Downstream consumers must therefore tolerate duplicates, for example by carrying an event id header and deduplicating in the query layer. The guarantee also only holds for sources that can push back. The exec source, which runs a command such as tail -F, cannot tell the process producing the data that a put failed, and Flume's own guide warns that such asynchronous sources offer no delivery guarantee. Prefer the taildir source, which records its read offsets in a position file and resumes after a restart.
Choosing a channel
The channel decides what survives a crash and how fast an agent can absorb a burst. The defaults below are from the 1.11.0 user guide; production values are usually much larger.
| Channel | Durability | Key settings (default) | Use when |
|---|---|---|---|
| memory | Lost on agent or host crash | capacity (100), transactionCapacity (100), byteCapacity | Losing the in-flight buffer is acceptable, or the source can replay |
| file | Write-ahead log on disk survives restarts | checkpointDir, dataDirs, capacity (1,000,000), transactionCapacity (10,000), checkpointInterval (30,000 ms), minimumRequiredSpace (524,288,000 bytes) | The default choice for log ingestion |
| kafka | Replicated in a Kafka topic | kafka.bootstrap.servers, kafka.topic, parseAsFlumeEvent (true) | You already run Kafka and want a shared, replayable buffer |
Put the file channel's dataDirs on different disks from the application's busy volumes; multiple data directories on separate disks increase throughput. The channel stops accepting puts and takes when free space drops below minimumRequiredSpace, which is better than corrupting its log but means a full disk becomes back-pressure on the source. useDualCheckpoints with a backupCheckpointDir shortens recovery after a crash that damages the checkpoint, because otherwise the channel replays its whole log to rebuild state. The spillable memory channel, which overflows from memory to a file channel, exists but the guide does not recommend it for production.
A two-tier topology, worked through
Why two tiers? If every application host wrote to HDFS directly, hundreds of writers would each hold open files and NameNode leases, producing many small files. Collectors fan in: a handful of agents write large files, and application hosts only need a local disk buffer. The tier-1 agent on each host:
a1.sources = r1
a1.channels = c1
a1.sinks = k1 k2
a1.sinkgroups = g1
a1.sources.r1.type = TAILDIR
a1.sources.r1.positionFile = /var/lib/flume/taildir_position.json
a1.sources.r1.filegroups = f1
a1.sources.r1.filegroups.f1 = /var/log/checkout/.*\.log
a1.sources.r1.headers.f1.service = checkout
a1.sources.r1.batchSize = 500
a1.sources.r1.interceptors = ts host
a1.sources.r1.interceptors.ts.type = timestamp
a1.sources.r1.interceptors.host.type = host
a1.sources.r1.channels = c1
a1.channels.c1.type = file
a1.channels.c1.checkpointDir = /data1/flume/checkpoint
a1.channels.c1.dataDirs = /data1/flume/data,/data2/flume/data
a1.channels.c1.capacity = 2000000
a1.sinks.k1.type = avro
a1.sinks.k1.hostname = collector-a.internal
a1.sinks.k1.port = 4545
a1.sinks.k1.batch-size = 500
a1.sinks.k1.channel = c1
a1.sinks.k2.type = avro
a1.sinks.k2.hostname = collector-b.internal
a1.sinks.k2.port = 4545
a1.sinks.k2.batch-size = 500
a1.sinks.k2.channel = c1
a1.sinkgroups.g1.sinks = k1 k2
a1.sinkgroups.g1.processor.type = failover
a1.sinkgroups.g1.processor.priority.k1 = 10
a1.sinkgroups.g1.processor.priority.k2 = 5
a1.sinkgroups.g1.processor.maxpenalty = 30000The timestamp interceptor stamps each event at ingestion, which the HDFS sink uses to choose the hour directory. If you need event time rather than arrival time, parse it from the log line with a regex extractor instead. The failover processor sends everything to the highest-priority healthy sink and retries a failed one after a penalty of up to maxpenalty milliseconds. Swap it for processor.type = load_balance with processor.backoff = true to spread load across both collectors.
The HDFS sink and the small-files trap
The collector writes to HDFS. Its most important settings are the roll triggers, and their defaults are a trap: hdfs.rollInterval is 30 seconds, hdfs.rollSize is 1,024 bytes and hdfs.rollCount is 10 events. Whichever fires first closes the file, so an unconfigured sink produces files of about one kilobyte. Every one costs NameNode memory and every query opens them all, which is exactly the problem described in the HDFS small files problem. Set the count and size triggers deliberately:
a2.sources = r1
a2.channels = c1
a2.sinks = k1
a2.sources.r1.type = avro
a2.sources.r1.bind = 0.0.0.0
a2.sources.r1.port = 4545
a2.sources.r1.channels = c1
a2.channels.c1.type = file
a2.channels.c1.checkpointDir = /data1/flume/checkpoint
a2.channels.c1.dataDirs = /data1/flume/data,/data2/flume/data,/data3/flume/data
a2.channels.c1.capacity = 2000000
a2.sinks.k1.type = hdfs
a2.sinks.k1.channel = c1
a2.sinks.k1.hdfs.path = hdfs://nameservice1/flume/%{service}/dt=%Y-%m-%d/hr=%H
a2.sinks.k1.hdfs.filePrefix = events-collector-a
a2.sinks.k1.hdfs.fileType = DataStream
a2.sinks.k1.hdfs.writeFormat = Text
a2.sinks.k1.hdfs.rollInterval = 3600
a2.sinks.k1.hdfs.rollSize = 125829120
a2.sinks.k1.hdfs.rollCount = 0
a2.sinks.k1.hdfs.idleTimeout = 600
a2.sinks.k1.hdfs.batchSize = 1000
a2.sinks.k1.hdfs.kerberosPrincipal = flume/collector-a.internal@EXAMPLE.COM
a2.sinks.k1.hdfs.kerberosKeytab = /etc/security/keytabs/flume.keytabA roll size of 120 MiB, a little under a 128 MiB block, gives files that fill roughly one block; a count of 0 disables the count trigger; an hourly interval and a 10-minute idle timeout close files for quiet services. DataStream writes plain bodies rather than the default SequenceFile. Files being written carry the .tmp suffix (hdfs.inUseSuffix) and are renamed on close, so downstream jobs must ignore in-progress files or they will read half-written data. A Hive table over the dt and hr directories then needs partitions added as hours close; Hive architecture explains the metastore side. On a secure cluster the principal and keytab authenticate the sink; see Hadoop Kerberos.
One collector writing large files is limited by the HDFS write pipeline it feeds, and hdfs.callTimeout (10,000 ms by default) bounds each HDFS call. A slow DataNode in the pipeline shows up as call timeouts, rolled-back batches and a growing channel, not as an error in the application.
Running and monitoring an agent
bin/flume-ng agent --conf conf --conf-file conf/collector.conf --name a2 \
-Dflume.monitoring.type=http -Dflume.monitoring.port=34545
curl -s http://collector-a.internal:34545/metrics
# {"CHANNEL.c1": {"ChannelSize": "233428", "ChannelCapacity": "2000000",
# "EventPutSuccessCount": "...", "EventTakeSuccessCount": "...", ...},
# "SINK.k1": {"EventDrainSuccessCount": "...", "ConnectionFailedCount": "...", ...}}The single most useful number is ChannelSize against ChannelCapacity for every channel. A channel that stays near empty is healthy. One that grows steadily means the sink drains slower than the source fills, and you have hours, not days, before the channel fills and the source starts failing puts. Alert on the ratio and on its slope. Compare EventPutSuccessCount with EventTakeSuccessCount over time to see the gap, and watch ConnectionFailedCount on Avro and HDFS sinks for downstream trouble.
Failure modes, trade-offs and migration
| Symptom | Cause | Fix |
|---|---|---|
| Millions of KB-sized files | Default roll settings | Set rollSize, rollCount = 0, rollInterval, idleTimeout |
| Channel full, source errors | Sink slower than source: slow HDFS, dead collector | Find the slow hop; add collectors; raise sink batch size |
| Data lost after a crash | Memory channel or exec source | File channel and taildir source |
| Duplicate rows downstream | Retry after a commit was lost | Expected with at-least-once; deduplicate by an id header |
| Readers see partial files | Jobs read .tmp files being written | Filter the in-use suffix or prefix |
| Events land in the wrong hour | Arrival-time timestamps, or local time on the sink | Extract event time; keep hosts in UTC |
| Agent out of memory | Huge memory channel, or many open HDFS files | Use a file channel; lower hdfs.maxOpenFiles (5,000 by default) |
Flume's strengths are simplicity, a local durable buffer on every host and direct HDFS writing with time partitioning. Its limits are that it is push-only, so consumers cannot replay data after it leaves a channel, and its connector and community activity has slowed. Kafka with Kafka Connect gives replay, many consumers per stream and a larger ecosystem, at the cost of running a Kafka cluster; Kafka streaming architecture covers it. A common migration keeps tier-1 Flume agents, switches their sink to Kafka, and moves HDFS writing to a consumer.
What to do next
- Inventory every agent and replace exec sources with taildir and memory channels with file channels where loss matters.
- Check every HDFS sink's roll settings, then count files per partition directory to confirm the fix.
- Enable HTTP monitoring and alert on ChannelSize as a fraction of ChannelCapacity and on its growth rate.
- Put file channel data directories on disks separate from application data, and enable dual checkpoints.
- Make downstream readers ignore in-progress .tmp files and tolerate duplicates.
- Add a second collector and a failover or load-balancing sink group so one collector can be restarted.
- Decide whether new streams should go to Kafka instead, and plan the tier-1 sink switch if so.