广告实时链路的地基,是先把散落在成百上千台机器上的曝光、点击、转化日志可靠地搬到 Kafka 上,再交给 Flink 做实时计费、归因与反作弊(这条链路的全景见 为什么 AdTech 偏爱 Kafka)。而”把日志从磁盘/网络搬到 Kafka/HDFS”这一步,最经典的组件就是 Apache Flume。这一篇把 Flume 从是什么、架构长什么样、怎么配、怎么用、优缺点在哪、什么时候该选它一路讲透。
一句话定位:Flume 是一个”可靠、可扩展、配置驱动”的分布式日志采集与聚合系统——核心抽象是
Source → Channel → Sink三段式的 Agent,用一段带事务的内部缓冲,把海量日志从产生端搬运到 Kafka / HDFS / 下一级 Agent,保证 at-least-once(不丢,可能重)。
TL;DR
- Flume = 配置化的搬运工:不写代码、只写
.conf,就能把日志从文件/网络采集、缓冲、发到目的地。核心是 Agent,一个 Agent =Source(进)→ Channel(存)→ Sink(出)。 - Event 是最小数据单元:
header(K-V 元数据)+ body(字节数组)。Source 把原始日志封装成 Event,Channel 缓冲 Event,Sink 消费 Event。 - 两段事务保不丢:
put 事务(Source → Channel)与take 事务(Channel → Sink)各自”成功才提交、失败就回滚”。Sink 发送失败时 Event 留在 Channel 重发——因此是 at-least-once(不丢但可能重),下游要靠幂等去重。 - Channel 是可靠性的开关:
Memory Channel快但宕机丢数据;File Channel落盘、崩溃可恢复(计费日志首选);Kafka Channel直接拿 Kafka 当 Channel、省掉一次 Sink。 - 三个高级组件撑起复杂拓扑:
Interceptor(在 Source 侧给 Event 打标/过滤/脱敏)、Channel Selector(一份数据复制或按 header 分流到多个 Channel)、Sink Processor(多 Sink 做 failover 或负载均衡)。 - 典型拓扑:采集层每台机器一个 Agent(
TailDir监控日志文件)→ 用AvroRPC 汇聚到少数几个汇聚 Agent → 汇聚 Agent 用Kafka Sink写进 Kafka。汇聚层把海量直连收敛成少数连接。 - 优点:零代码配置化、和 Hadoop/Kafka 生态天生契合、可靠性机制完整、拓扑灵活。缺点:吞吐/资源不如 Vector/Fluent Bit、清洗能力弱、社区已进入维护状态、Memory Channel 易丢、HDFS Sink 易产生小文件。
- 选型:存量 Hadoop 平台、要写 HDFS/Kafka → Flume 仍是顺手之选;新建 / 云原生 / K8s → 优先 Fluent Bit、Vector;需要重清洗 → Logstash / Fluentd。广告核心计费链路的判断是”能不能丢”——不能丢就
File Channel + acks=-1。
Table of contents
Open Table of contents
1. Flume 是什么,解决什么问题
在没有专门采集组件的年代,“把日志弄到大数据平台”通常是一堆脆弱的脚本:crontab + scp 定时拷文件、或自己写程序 tail 文件再往下游发。这些方案有共同的痛点:断点续传难做(进程重启就重复或漏读)、宕机会丢数据、加一个新目的地就要改代码、上千台机器的运维不可控。
Flume 就是为解决这类”日志从产生端到集中存储的可靠搬运”而生的。它的设计取舍很清晰:
- 配置驱动、零代码:绝大多数采集需求只写一个
.conf文件就能跑起来,不用编译打包。 - 可靠传输:内建事务 + 可落盘的 Channel,保证 at-least-once——正常情况下数据不丢。
- 和大数据生态天生一家:Flume 是 Hadoop 生态原生项目,写 HDFS / HBase / Kafka 是一等公民,配置项开箱即用。
- 可组合的拓扑:Agent 之间能串联、汇聚(fan-in)、扇出(fan-out),拼出采集层 → 汇聚层的多级架构。
它不是什么:Flume 不是流计算引擎(不做复杂 join/聚合,那是 Flink 的活),也不是消息队列(缓冲能力有限,削峰解耦要靠 Kafka)。它只做好一件事——可靠地把日志搬过去。
2. 核心概念与架构:Agent / Event / Source / Channel / Sink
Flume 的运行单元是 Agent——一个独立的 JVM 进程。每个 Agent 内部由三大组件串成一条流水线:Source(数据从哪来)→ Channel(中间缓冲)→ Sink(数据到哪去)。
Flume Agent 的三段式与两段事务:Source 把外部日志封装成 Event 写入 Channel(put 事务),Sink 从 Channel 取出 Event 发往目的地(take 事务)。两段事务都”成功才提交、失败就回滚”,这正是 at-least-once(不丢但可能重)的来源。
2.1 Event:Flume 里流动的最小数据单元
Flume 内部搬运的一切都是 Event,结构极简:
Event = { headers: Map<String, String>, // 元数据:时间戳、来源、日志类型、路由 key…
body: byte[] } // 真正的日志内容(通常是一行原始日志)
body 是原始日志的字节;headers 是可选的键值对元数据,拦截器(§3.1)可以往里塞标记,Channel Selector(§3.2)可以按 header 决定数据往哪走。理解”一切皆 Event、header 可路由”是读懂 Flume 高级玩法的钥匙。
2.2 Source:数据从哪来
Source 负责接收数据、封装成 Event 写入 Channel。常用的几类:
| Source | 用途 | 广告场景 |
|---|---|---|
| TailDir | 监控一批日志文件/目录,支持断点续传(positionFile 记录读到哪) | ⭐ 采集本地滚动日志的首选 |
| Avro / Thrift | 接收上游 Agent 通过 RPC 发来的数据 | ⭐ 汇聚层收采集层的数据 |
| Kafka Source | 从 Kafka topic 消费 | 从 Kafka 再落 HDFS |
| Exec | 执行命令(如 tail -F)取输出 | 不推荐,进程挂了就丢、无断点续传 |
| Spooling Directory | 监控一个目录,处理”完成的文件” | 批量文件采集 |
老教程里常见的
Exec Source + tail -F在生产要避免:它不记录读取位置,Agent 重启会重复或丢数据。要采集滚动日志文件,用 TailDir Source——它把每个文件的读取偏移量持久化到positionFile,重启后接着读。
2.3 Channel:中间缓冲,也是可靠性的总开关
Channel 是 Source 和 Sink 之间的缓冲队列,Flume 的可靠性等级基本由选哪种 Channel 决定:
| Channel | 特点 | 广告场景 |
|---|---|---|
| Memory Channel | Event 存内存队列,快、吞吐高;但进程崩溃 / 断电 → 队列数据全丢 | 能容忍少量丢失的场景(如监控指标) |
| File Channel | Event 落本地磁盘(含 WAL checkpoint),慢一些;Agent 崩溃重启后可从磁盘恢复未发送的 Event | ⭐ 广告计费这类”不能丢”的日志首选 |
| Kafka Channel | 直接把 Kafka 当作 Channel(相当于 Source → Kafka),可省掉 Sink;兼具 Kafka 的持久化 + 副本可靠性,又少一跳 | ⭐ 现代常用,直连 Kafka 又省一跳 |
一句话记忆:要快选 Memory、要稳选 File、要直连 Kafka 又省一跳选 Kafka Channel。 广告日志的默认答案是 File Channel(或 Kafka Channel)。
2.4 Sink:数据到哪去
Sink 从 Channel 取出 Event 发往目的地。常用:
| Sink | 用途 | 广告场景 |
|---|---|---|
| Kafka Sink | 写入 Kafka topic | ⭐ 广告链路最常用,下游接 Flink |
| HDFS Sink | 写入 HDFS,按时间分区落盘 | 离线数仓 T+1 |
| Avro / Thrift Sink | 通过 RPC 发给下一级 Agent | 多级汇聚(采集层 → 汇聚层) |
| HBase / ES Sink | 写入 HBase / Elasticsearch | 实时点查 / 检索 |
2.5 可靠性:两段事务与 at-least-once
Flume “不丢”的秘密是 Channel 两端各有一个事务(见上图):
- put 事务(Source → Channel):Source 把一批 Event 批量放进 Channel,全部成功才 commit;任何一条失败就 rollback,这批 Event 不进 Channel(Source 会重新采集)。
- take 事务(Channel → Sink):Sink 从 Channel 取一批 Event 发往目的地,目的地确认接收成功才 commit(此时 Event 才从 Channel 真正删除);发送失败则 rollback,Event 留在 Channel 里等下次重发。
这套机制的结果是 at-least-once:数据不会丢,但”发送成功了、commit 前 Agent 崩溃”这类边界会导致重复。所以广告计费链路的正确姿势是——Flume 保证不丢(File Channel + 事务),下游 Flink/消费端用业务唯一键幂等去重,二者配合等价于端到端”有效 exactly-once”(这与 Kafka 计费对账 的思路一致)。
3. 高级组件:拦截器、选择器与 Sink 处理器
三段式只是骨架。真正让 Flume 能拼出复杂采集拓扑的,是三个可插拔的高级组件。
3.1 Interceptor:在 Source 侧加工 Event
拦截器挂在 Source 之后、写入 Channel 之前,可以链式串多个,对每个 Event 做轻量加工:
- Timestamp / Host Interceptor:给 header 加时间戳、主机名(HDFS Sink 按时间分区就靠它)。
- Static Interceptor:给 header 塞固定 K-V(如
logtype=impression),供下游路由。 - Regex Filtering Interceptor:按正则过滤或保留日志行(丢掉心跳、脏行)。
- 自定义 Interceptor:实现
Interceptor接口做脱敏、字段抽取等(但重清洗别放这里,会拖慢采集,交给下游 Flink)。
3.2 Channel Selector:一份数据分流到多个 Channel
当一个 Source 要把数据发到多个 Channel(进而多个 Sink)时,由 Channel Selector 决定怎么分:
- Replicating(复制,默认):每条 Event 复制进所有 Channel——比如同一份日志一路进 Kafka(实时)、一路进 HDFS(离线)。
- Multiplexing(多路复用):按 Event header 的值路由到不同 Channel——比如按
header.logtype把曝光日志和点击日志分到不同 Channel、发往不同 topic。
3.3 Sink Processor / Sink Group:多 Sink 的容错与均衡
把多个 Sink 组成一个 Sink Group,由 Sink Processor 决定行为:
- Failover(故障转移):Sink 按优先级排序,主 Sink 挂了自动切到备用 Sink——目的地高可用。
- Load Balancing(负载均衡):多个 Sink 轮询 / 随机分担,把一个 Channel 的数据分摊出去,提升吞吐、打散下游压力。
这三个组件的组合能力,是 Flume 相较轻量 agent(Filebeat 等)的差异化优势:在采集层就能完成打标、过滤、分流、多目的地容错,而不必所有逻辑都堆到下游。
4. 配置与使用实战
Flume 的使用 = 写 .conf + 启动 Agent。配置文件的通用范式是先声明组件名,再逐个配置属性,最后把 Source/Sink 绑到 Channel 上。
4.1 采集本地日志 → Kafka(TailDir + File Channel + Kafka Sink)
这是广告采集层最典型的一个 Agent:监控日志目录、File Channel 保不丢、写进 Kafka。
# ad-collect-agent.conf
# ① 声明组件
a1.sources = r1
a1.channels = c1
a1.sinks = k1
# ② Source:TailDir 监控多组日志目录,断点续传
a1.sources.r1.type = TAILDIR
a1.sources.r1.channels = c1
a1.sources.r1.positionFile = /var/flume/taildir_position.json
a1.sources.r1.filegroups = imp clk
a1.sources.r1.filegroups.imp = /data/logs/ad/impression/.*log
a1.sources.r1.filegroups.clk = /data/logs/ad/click/.*log
# 用拦截器给不同日志打上 logtype,供下游区分
a1.sources.r1.interceptors = i1
a1.sources.r1.interceptors.i1.type = static
a1.sources.r1.interceptors.i1.key = logtype
a1.sources.r1.interceptors.i1.value = ad
# ③ Channel:File Channel 落盘,Agent 崩溃可恢复(计费日志关键)
a1.channels.c1.type = file
a1.channels.c1.checkpointDir = /var/flume/checkpoint
a1.channels.c1.dataDirs = /var/flume/data
a1.channels.c1.capacity = 1000000
a1.channels.c1.transactionCapacity = 10000
# ④ Sink:写入 Kafka,acks=-1 等 ISR 副本全部确认
a1.sinks.k1.type = org.apache.flume.sink.kafka.KafkaSink
a1.sinks.k1.channel = c1
a1.sinks.k1.kafka.bootstrap.servers = kafka1:9092,kafka2:9092,kafka3:9092
a1.sinks.k1.kafka.topic = ad_raw_log
a1.sinks.k1.kafka.producer.acks = -1
a1.sinks.k1.kafka.producer.compression.type = lz4
a1.sinks.k1.kafka.flumeBatchSize = 2000
启动:
flume-ng agent \
--name a1 \
--conf /etc/flume/conf \
--conf-file /etc/flume/conf/ad-collect-agent.conf \
-Dflume.root.logger=INFO,console
4.2 多级汇聚:采集层 → 汇聚层(Avro Sink → Avro Source)
采集层不直接连 Kafka,而是通过 Avro RPC 把数据发给汇聚 Agent——这就是扇入(fan-in):把多台机器的数据收敛到少数汇聚 Agent(好处见 §5)。
采集端(Sink 换成 Avro,指向汇聚 Agent):
# 采集端 sink:发给汇聚 Agent
a1.sinks.k1.type = avro
a1.sinks.k1.channel = c1
a1.sinks.k1.hostname = collector-agent.internal
a1.sinks.k1.port = 4545
a1.sinks.k1.batch-size = 2000
汇聚端(Source 用 Avro 接收,再统一写 Kafka):
# 汇聚端 source:接收采集端的 Avro RPC
c1.sources.r1.type = avro
c1.sources.r1.channels = ch1
c1.sources.r1.bind = 0.0.0.0
c1.sources.r1.port = 4545
# ch1 用 File Channel,sink 再写 Kafka(同 §4.1)
4.3 落 HDFS:按小时分区、避免小文件
离线链路常用 HDFS Sink,配置的重点是滚动策略——控制多久/多大生成一个文件,避免 HDFS 小文件泛滥(§8 会展开)。
a1.sinks.k2.type = hdfs
a1.sinks.k2.channel = c1
a1.sinks.k2.hdfs.path = hdfs://nn/ad/raw/dt=%Y%m%d/hour=%H
a1.sinks.k2.hdfs.filePrefix = ad
a1.sinks.k2.hdfs.fileType = DataStream
# 滚动策略:满足任一条件就切新文件——按时间/大小滚,别按条数
a1.sinks.k2.hdfs.rollInterval = 3600 # 每小时滚一次
a1.sinks.k2.hdfs.rollSize = 134217728 # 或达到 128MB 滚
a1.sinks.k2.hdfs.rollCount = 0 # 关掉按条数滚(否则文件极碎)
a1.sinks.k2.hdfs.round = true
a1.sinks.k2.hdfs.roundValue = 1
a1.sinks.k2.hdfs.roundUnit = hour
5. 广告场景:一条不丢的日志采集链路
把前面的组件拼起来,就是广告业界最常见的采集层 → 汇聚层 → Kafka → 双消费架构。
广告日志采集的经典分层:采集层每机一个 TailDir Agent,经 Avro 汇聚到少数汇聚 Agent,再由 Kafka Sink 统一写入 Kafka,下游 Flink(实时)与 HDFS(离线)各起一个消费组。汇聚层把海量直连收敛成少数连接,也是统一压缩/限流/加密的地方。
这条链路是两个动作的叠加:采集层多 Agent **扇入(fan-in)汇聚层收敛连接,Kafka 再扇出(fan-out)**给 Flink(实时)与 HDFS(离线)两个消费组。
为什么要中间加一层汇聚,而不是每台机器直连 Kafka?
- 收敛连接数:上千台机器直连 Kafka = 上千个 producer 连接压在 broker 上;经汇聚层收敛成几十个连接,broker 压力小得多。
- 统一治理点:压缩、限流、加密、打标都在汇聚层统一做,采集层保持极简。
- 削峰的一层缓冲:汇聚 Agent 的 File Channel 又是一层磁盘缓冲,Kafka 抖动时能兜一会儿。
这条链路”不丢”的三个关键点:
- 采集层用 File Channel——Agent 崩溃重启不丢未发送的 Event;
TailDir positionFile保证不重读不漏读。 - Kafka Sink
acks=-1+ topic 3 副本 +min.insync.replicas=2——写入侧不丢(原理见 Kafka 可靠性与 Exactly-Once)。 - 下游 Flink 用业务唯一键幂等去重——消化 Flume at-least-once 带来的重复,端到端等价”有效 exactly-once”。
“日志怎么进管道”其实是一条谱系,不止两种:① 应用直写 Kafka——内嵌 Producer、不落盘也不经 Flume,延迟最低,但可靠性全压在 Kafka 自己身上;② 落盘 + TailDir/Flume 采集——最抗丢(Agent 崩了磁盘还在),曝光/计费这类”绝不能丢”的日志常走,代价是延迟到秒级;③ 应用经 RPC 把事件推给 Flume——客户端内存队列 + Agent Memory Channel 压低延迟,再由 Sink 写 Kafka,既不落盘也不直连 Kafka,却仍能吃到 Flume 的分流 / 双写 / failover。下面 §6 的真实 DSP 走的正是第 ③ 条——连毫秒级的竞价事件也经 Flume(而非直写 Kafka),靠”内存队列 + 全 Memory Channel +
acks=1”把延迟压下来。很多广告公司几条链路并存。
6. 真实生产案例:一份 DSP 汇聚 Agent 的配置拆解
前面的组件在真实 DSP 里怎么拼?下面拆解一份生产环境的汇聚 Agent 配置(已脱敏)——从上游 Bidder 侧的客户端封装(6.1),到 Agent 端的 multiplexing 分流、双写、failover sink group,把 §3 的组件全用上了,是一个”教科书组件 → 真实系统”的好样本。
说明:本节做了脱敏处理——Kafka 地址、OSS 凭证、镜像仓库、资源组 ID 等一律用占位符;事件类型名、topic、内部路径、Channel / Sink 组件名等内部标识符已改写为示意值(如用
event/bill代表单路 / 双路两类事件),客户端类名(FlumeRpcClient等)为通用示意命名。重点在架构与取舍,不在具体名字。真实系统里请勿把凭证明文写进配置或部署模板。
6.1 上游 Bidder 侧:FlumeRpcClient 如何把事件推进来
上面的 Avro Source 收的是上游 Bidder(竞价引擎)推来的事件。Bidder 里封装了一套轻量客户端,把「业务代码产日志」和「网络发往 Flume」解耦成 生产者 → 内存队列 → 消费者 → RPC 四段,核心目的是让竞价主链路零阻塞:
Bidder 侧客户端四段式:业务线程只管把日志 offer 进内存队列(零阻塞),单线程消费者按 msgType 分类攒批,再由 FlumeRpcClient 负载均衡地 Avro RPC 打到汇聚 Agent 的 Source。
① FlumeEventProducer(生产者):业务侧调各类 sendXxxMessage,把日志封成 Flume Event(body = UTF-8 日志正文,headers 带 msgType / msgVersion / timestamp),非阻塞地 offer 进一个内存队列,超时 20ms 入队失败就记告警。关键在 msgType 的拼法——类型 + 版本,如 event + 1_0 → event1_0、bill1_0:
// msgType = 类型 + 版本号(1.0 → 1_0):event1_0 / bill1_0
headers.put("msgType", msgType + version.replace(".", "_"));
// 封成 Flume Event 后非阻塞入队,20ms 入不进就告警(客户端第一道背压/丢弃点)
Event event = EventBuilder.withBody(data, StandardCharsets.UTF_8, headers);
boolean added = FlumeRpcClient.eventQueue.offer(event, 20, TimeUnit.MILLISECONDS);
if (!added) log.error("Flume event queue is full ...");
这串
msgType值,正是 Flume Agent 端 multiplexing 选择器的路由 key(§6.2 的selector.mapping.<msgType>)——客户端设 header、服务端按 header 分流,两端靠同一个msgType串起来。
② LinkedBlockingQueue(解耦缓冲):生产与发送彻底解耦。业务(竞价)线程只管入队、绝不阻塞在网络 IO 上;队列满则丢弃并告警——用”可能丢一点日志”换”竞价主链路不被拖慢”。
③ FlumeEventConsumer(消费者,单线程 daemon):一个后台线程不断 drainTo(batch=10),按 msgType 把事件分类攒批(计费类事件立即发、普通事件攒够 10 条再发),交给 RpcClient。注意——这里的”分类”只是客户端按类型分批,真正”写进不同 Channel”的动作发生在 §6.2 的 Agent multiplexing 选择器,不在客户端。
④ FlumeRpcClient(RPC 发送):包住 Flume 官方 RpcClient,用 default_loadbalance + 轮询(round-robin)在多台汇聚 Agent 间负载均衡,带 backoff、失败自动重连重试、批量 appendBatch。它对应的正是 §6.2 那个 :4141 的 Avro Source。
客户端这条链路的取舍和 Agent 侧一脉相承:内存队列 + 单线程批量 + 负载均衡 RPC,业务线程零阻塞、发送侧攒批提吞吐、多 Agent 容错;代价是队列满会丢、且这段内存队列同样”进程被杀就丢”(和 §6.5 Agent 侧 Memory Channel 的取舍如出一辙)。
再往外拉一层看:“采集端”其实有一大把。 上面画的是一个 Bidder 实例内部的四段流水;而竞价引擎是横向扩容的——线上跑着几十到上百个 Bidder 实例,每个实例都内嵌了一份 FlumeRpcClient。它们彼此独立、都把事件通过 Avro RPC 推给那少数几台汇聚 Agent,这就是 §5 说的**扇入(fan-in)**在这套系统里的真身:
- “多个采集 Agent”= 多个 Bidder 实例里的客户端:Bidder 横向扩多少个实例,就有多少个”采集端”在往 Flume 推;每个只管自己进程内的事件、彼此独立。
- 靠
FlumeRpcClient的 round-robin 负载均衡扇入:上百个客户端不各自直连 Kafka(那会把 broker 连接数打爆),而是先把连接收敛到少数几台汇聚 Agent。 - 一个 Avro Source 承接全部连接:§6.2 那个
:4141、16 线程的 Source,就是用来统一承接这上百个 Bidder 客户端的——这正是”扇入收敛连接数”的价值所在。
6.2 单 Avro Source + multiplexing:按 msgType 分流到十几路
整个 Agent 只有一个 Avro Source(:4141,16 线程,收上游 Bidder 通过 FlumeRpcClient 推来的事件),用 multiplexing selector 按 Event 的 msgType header 把十几类事件路由到各自的 Channel:
agent.sources.avroSrc.type = avro
agent.sources.avroSrc.bind = 0.0.0.0
agent.sources.avroSrc.port = 4141
agent.sources.avroSrc.threads = 16
# 按 header msgType 多路分流(§3.2 Multiplexing 的实战)
agent.sources.avroSrc.selector.type = multiplexing
agent.sources.avroSrc.selector.header = msgType
# 大多数事件是「单路」:一种 msgType → 一个 Channel。以普通事件 event 为例:
agent.sources.avroSrc.selector.mapping.event1_0 = memChannelEvent
# 少数货币化事件是「双路」:一种 msgType → 两个 Channel(复制双写,见 6.3)。以计费事件 bill 为例:
agent.sources.avroSrc.selector.mapping.bill1_0 = memChannelBill memChannelKafkaBill
# … 其余出价、曝光、点击、转化、模型打分等十几类事件 → 各自的 Channel,各按上面两种模式之一映射
映射其实只有两种模式:
- 单路(大多数事件,如普通事件
event):一种msgType→ 一个 Channel → 一个 Sink,走一条链路(通常是落文件 → OSS)。 - 双路(货币化事件,如计费事件
bill):一种msgType→ 两个 Channel,复制成两份分别走离线与实时(见 6.3)。
一句话:竞价链路的十几种事件(出价、竞得、曝光、点击、转化、模型打分……)全从一个源头进来,按类型分流,再按”单路 / 双路”决定后续命运。
6.3 双写:货币化事件同时进「文件(→OSS)」和「Kafka(→实时)」
注意上面 bill1_0 被映射到两个 Channel——这是 §3.2 的 Replicating 效果(一份 Event 复制两份),一份走离线、一份走实时:
一份 bill 事件的两条命运:离线路经自研 FileSink 落本地大文件、由脚本上传 OSS 供数仓补数/对账;实时路经 failover 组写 Kafka 给 Flink 做实时计费。全程 Memory Channel,靠双写 + failover 补可靠性。
# —— 离线路:FileSink 落本地文件,再由脚本上传 OSS ——
agent.sinks.billFileSink.type = com.example.flume.sink.FileSink # 自研 FileSink
agent.sinks.billFileSink.channel = memChannelBill
agent.sinks.billFileSink.file.path = /data/logs/dsp
agent.sinks.billFileSink.file.rollSize = 536870912 # 512MB 滚一个文件
agent.sinks.billFileSink.file.maxWritingTimeMin = 10
agent.sinks.billFileSink.file.script = /opt/scripts/upload-oss.sh # 上传 OSS
# —— 实时路:写 Kafka,下游 Flink 消费 ——
agent.sinks.billKafkaSink.type = org.apache.flume.sink.kafka.KafkaSink
agent.sinks.billKafkaSink.channel = memChannelKafkaBill
agent.sinks.billKafkaSink.kafka.bootstrap.servers = <kafka-broker-1>:9092,<kafka-broker-2>:9092,<kafka-broker-3>:9092
agent.sinks.billKafkaSink.kafka.topic = ad-billing
agent.sinks.billKafkaSink.kafka.producer.acks = 1 # 注意:acks=1 而非 -1(延迟优先)
agent.sinks.billKafkaSink.kafka.producer.linger.ms = 100
agent.sinks.billKafkaSink.kafka.producer.compression.type = snappy
agent.sinks.billKafkaSink.kafka.flumeBatchSize = 200
一份数据两条命运:离线路攒成大文件传 OSS,供数仓补数/对账;实时路进 Kafka 给 Flink 做实时计费/归因。多个货币化事件类型还共用同一个 topic
ad-billing,靠字段在下游区分。
6.4 failover sink group:Kafka 挂了自动降级写本地盘
实时路的 Kafka Sink 不是单打独斗,而是和一个自研 FileSink 组成 failover 组——Kafka 不可用时自动切到本地文件 spill,事后再补:
agent.sinkgroups.billGroup.sinks = billKafkaSink billFailoverSink
agent.sinkgroups.billGroup.processor.type = failover
agent.sinkgroups.billGroup.processor.priority.billKafkaSink = 10 # 主:写 Kafka
agent.sinkgroups.billGroup.processor.priority.billFailoverSink = 5 # 备:Kafka 挂了落盘
agent.sinkgroups.billGroup.processor.maxpenalty = 10000 # 故障 sink 最长退避 10s
这正是 §3.3 Failover Sink Processor 的实战:几条货币化链路(如 bill)各配一个 failover 组,把”Kafka 抖动 = 丢计费事件”的风险兜住。
6.5 可靠性取舍:为什么这里敢全用 Memory Channel
有意思的是——这份配置所有 Channel 都是 memory(离线文件路 capacity=10 万 / transactionCapacity=1 万,Kafka 路更小、约 1 万 / 1000),和 §2.3 建议的”计费日志用 File Channel”正好相反。它的可靠性不靠 Channel,而是靠一套组合拳:
- 双写(6.3):离线文件路 + 实时 Kafka 路互为冗余;
- failover 落盘(6.4):Kafka 挂了 spill 到本地文件;
- 落盘到持久卷:文件路写在挂载的持久存储上,Agent 重建后数据还在。
这是一个**“延迟 / 吞吐优先”的取舍(配合 acks=1、snappy、小 batch、linger.ms=100)。但要清醒认识残留风险**:failover 只在”Kafka 不可用”时触发;如果 进程被硬 kill(OOMKill、节点宕机,跳过优雅停机),Memory Channel 里在途的事件会直接丢。能否接受,取决于这些事件有没有离线文件路兜底、以及对账能否补回。
6.6 从这份配置能学到什么
- ✅ 一个 Source 按 header 分流十几路 + 关键事件双写,是 DSP 采集汇聚层的成熟范式。
- ✅ failover sink group 把”Kafka 抖动丢数据”兜住,比”死等 Kafka”更稳。
- ⚠️ 全 Memory Channel 是延迟换来的,硬 kill 仍有丢数窗口——靠双写 + failover 落盘缓解,但要评估计费口径能否接受。
- ⚠️ 凭证绝不能明文进 env:务必改 Secret / KMS;
acks=1、rollInterval语义、多事件共用 topic 也都是评审时要确认的点。
7. 优缺点分析
综合来看,Flume 的取舍非常鲜明:
| 维度 | 优点 ✅ | 缺点 ⚠️ |
|---|---|---|
| 上手成本 | 配置驱动、零代码,写 .conf 即用 | 配置繁琐、组件语义要吃透(Channel/事务/滚动策略易配错) |
| 可靠性 | 事务 + File Channel,at-least-once 不丢 | 只有 at-least-once,会重复;Memory Channel 崩溃即丢 |
| 生态契合 | Hadoop 原生,写 HDFS/HBase/Kafka 一等公民 | 强绑大数据生态,云原生/K8s 场景不够顺手 |
| 拓扑能力 | 拦截器/选择器/Sink 处理器,采集层就能打标/分流/容错 | 复杂拓扑配置维护成本高 |
| 性能/资源 | 单 Agent 吞吐够用 | 基于 JVM,资源占用与吞吐不如 Vector / Fluent Bit |
| 数据加工 | 轻量拦截器够用 | 清洗/转换能力弱,重加工要下沉到 Flink |
| 运维/生态活跃度 | 稳定、久经考验 | 社区已进入维护状态,新特性少、新项目选它的越来越少 |
| HDFS 写入 | 原生支持、按时间分区 | 易产生小文件,滚动策略配不好会拖垮 NameNode |
一句话总结:Flume 的价值在”可靠 + 生态契合 + 配置化”,短板在”性能/资源、清洗能力、以及日渐式微的社区”。
8. 生产实践与踩坑
- 别用 Memory Channel 跑计费日志:Memory Channel 一旦 Agent OOM / 被 kill,队列里的 Event 全丢。凡是”不能丢”的日志一律 File Channel(或 Kafka Channel)。代价是磁盘 IO,吞吐会降。(§6 的生产案例给了另一种折中——“Memory Channel + 双写 + failover 落盘”,用吞吐换可靠性,但硬 kill 仍有丢数窗口,需按计费口径评估。)
- Channel 满了会反压:Sink 发送变慢(如 Kafka 抖动)时,Channel 会被填满,进而 Source 的 put 事务失败、采集停滞。要监控
ChannelSize/ChannelFillPercentage,把capacity调够,并给 File Channel 足够磁盘。 - HDFS 小文件是头号坑:
hdfs.rollCount默认按条数滚(默认 10 条一个文件!)会产生海量小文件压垮 NameNode。务必rollCount=0,改用rollInterval(按时间)+rollSize(按大小 ~128MB)滚动。 batchSize与transactionCapacity要匹配:Sink 的batchSize不能大于 Channel 的transactionCapacity,否则一批取不完会报错。批调大能提吞吐,但也放大重复窗口。- TailDir 的 positionFile 要落在可靠路径:
positionFile记录每个文件读到哪,若放在临时目录被清掉,重启会从头重读导致大量重复。 - 重复是常态,下游必须幂等:Flume 是 at-least-once,网络重传、commit 前崩溃都会造成重复。去重的责任在下游(Flink 状态去重 / 消费端业务唯一键),不要指望 Flume 帮你精确一次。
- Exec Source 只适合玩具:
tail -F式的 Exec Source 无断点续传、进程挂了就丢,生产采集文件一律用 TailDir。
9. 选型:Flume 还是 Filebeat / Fluentd / Vector?
Flume 只是日志采集 agent 里的一员。放到今天的技术格局里横向对比:
| 组件 | 语言 | 资源占用 | 特点 | 适用 |
|---|---|---|---|---|
| Flume | Java | 较高 | Hadoop 原生,写 HDFS/Kafka 一把梭,可靠性完整 | 存量大数据平台、写 HDFS/Kafka |
| Filebeat | Go | 低 | 轻量,专注采集文件转发,ELK 栈标配 | 每机部署采集本地日志 |
| Logstash | JRuby | 高 | 清洗/转换能力强、插件丰富 | 日志需复杂过滤解析 |
| Fluentd / Fluent Bit | Ruby+C / C | 中 / 极低 | CNCF 项目,云原生首选,Fluent Bit 极轻 | K8s DaemonSet 采集容器日志 |
| Vector | Rust | 极低 | 新一代,性能强、可观测性好 | 追求性能/现代化的新平台 |
怎么选(决策路径):
- 存量 Hadoop 平台、目的地是 HDFS/Kafka、团队熟 Flume → Flume 仍是顺手、稳妥的选择。
- Kubernetes / 云原生 / 容器日志 → Fluent Bit(DaemonSet)或 Fluentd,别硬塞 Flume。
- 追求极致性能 / 新建现代化采集 → Vector(Rust,性能与资源俱佳)。
- ELK 日志分析栈 → Filebeat 采集 + Logstash 清洗(可选)。
- 需要在采集端就做复杂清洗/解析 → Logstash / Fluentd / Vector,而不是 Flume。
现状与趋势:Flume 社区已进入维护状态(活跃度低、少有新特性),大量新项目转向 Vector / Fluent Bit。但在已建成的 Hadoop 数仓体系里,Flume 依然广泛存活——存量系统不会轻易替换一个”能稳定跑、和 HDFS/Kafka 无缝”的采集层。看你是”老平台”还是”新平台”,是选 Flume 的分水岭。
参考
- Apache Flume. Flume User Guide(官方用户指南):Source / Channel / Sink、拦截器 / 选择器 / Sink 处理器、全部配置项的一手权威文档。
- Apache Flume. Getting Started / Architecture:Agent 模型、Event 与可靠性机制的官方说明。
- Vector. Vector Documentation:现代高性能采集器,选型对比时的参照系。