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.

Advertisement

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.

Advertisement

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.

ChannelDurabilityKey settings (default)Use when
memoryLost on agent or host crashcapacity (100), transactionCapacity (100), byteCapacityLosing the in-flight buffer is acceptable, or the source can replay
fileWrite-ahead log on disk survives restartscheckpointDir, dataDirs, capacity (1,000,000), transactionCapacity (10,000), checkpointInterval (30,000 ms), minimumRequiredSpace (524,288,000 bytes)The default choice for log ingestion
kafkaReplicated in a Kafka topickafka.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

TIER 1: on each app hostTIER 2: collectorsTaildir sourcepositionFile, interceptorsFile channelWAL on local diskAvro sink k1priority 10Avro sink k2priority 5put txntakefailover sink group g1Collector A: Avro sourceFile channelHDFS sinkroll by size, DataStreamprimarystandbyCollector B (same shape)HDFS/flume/service/dt=/hr=close + renameEvery hop is a channel transaction: an event leavesa channel only after the next hop has accepted it.Result: at-least-once delivery, with duplicates on retry.
Tier-1 agents on application hosts tail logs into a durable file channel and forward over Avro to collectors, with failover. Collectors buffer again and write time-partitioned files to HDFS. Each hop commits only after the next hop accepted the batch.

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 = 30000

The 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.keytab

A 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

SymptomCauseFix
Millions of KB-sized filesDefault roll settingsSet rollSize, rollCount = 0, rollInterval, idleTimeout
Channel full, source errorsSink slower than source: slow HDFS, dead collectorFind the slow hop; add collectors; raise sink batch size
Data lost after a crashMemory channel or exec sourceFile channel and taildir source
Duplicate rows downstreamRetry after a commit was lostExpected with at-least-once; deduplicate by an id header
Readers see partial filesJobs read .tmp files being writtenFilter the in-use suffix or prefix
Events land in the wrong hourArrival-time timestamps, or local time on the sinkExtract event time; keep hosts in UTC
Agent out of memoryHuge memory channel, or many open HDFS filesUse 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

  1. Inventory every agent and replace exec sources with taildir and memory channels with file channels where loss matters.
  2. Check every HDFS sink's roll settings, then count files per partition directory to confirm the fix.
  3. Enable HTTP monitoring and alert on ChannelSize as a fraction of ChannelCapacity and on its growth rate.
  4. Put file channel data directories on disks separate from application data, and enable dual checkpoints.
  5. Make downstream readers ignore in-progress .tmp files and tolerate duplicates.
  6. Add a second collector and a failover or load-balancing sink group so one collector can be restarted.
  7. Decide whether new streams should go to Kafka instead, and plan the tier-1 sink switch if so.
Key takeaway: Flume moves events through sources, channels and sinks, and every hop is a transaction: events leave a channel only after the next hop accepts them, which gives at-least-once delivery with duplicates on retry. Durability comes from the file channel and replayable sources such as taildir, not from the memory channel or exec source. A two-tier layout with failover keeps HDFS writers few, and the HDFS sink's roll defaults must be overridden or it writes kilobyte files. Watch channel fill levels, expect duplicates, and weigh Kafka for new pipelines.