Skip to content
Charles Shao
Go back

Apache Flume:从 Source·Channel·Sink 到广告日志不丢链路与选型

views

广告实时链路的地基,是先把散落在成百上千台机器上的曝光、点击、转化日志可靠地搬到 Kafka,再交给 Flink 做实时计费、归因与反作弊(这条链路的全景见 为什么 AdTech 偏爱 Kafka)。而”把日志从磁盘/网络搬到 Kafka/HDFS”这一步,最经典的组件就是 Apache Flume。这一篇把 Flume 从是什么、架构长什么样、怎么配、怎么用、优缺点在哪、什么时候该选它一路讲透。

一句话定位:Flume 是一个”可靠、可扩展、配置驱动”的分布式日志采集与聚合系统——核心抽象是 Source → Channel → Sink 三段式的 Agent,用一段带事务的内部缓冲,把海量日志从产生端搬运到 Kafka / HDFS / 下一级 Agent,保证 at-least-once(不丢,可能重)

TL;DR

Table of contents

Open Table of contents

1. Flume 是什么,解决什么问题

在没有专门采集组件的年代,“把日志弄到大数据平台”通常是一堆脆弱的脚本:crontab + scp 定时拷文件、或自己写程序 tail 文件再往下游发。这些方案有共同的痛点:断点续传难做(进程重启就重复或漏读)、宕机会丢数据、加一个新目的地就要改代码、上千台机器的运维不可控

Flume 就是为解决这类”日志从产生端到集中存储的可靠搬运”而生的。它的设计取舍很清晰:

不是什么:Flume 不是流计算引擎(不做复杂 join/聚合,那是 Flink 的活),也不是消息队列(缓冲能力有限,削峰解耦要靠 Kafka)。它只做好一件事——可靠地把日志搬过去

2. 核心概念与架构:Agent / Event / Source / Channel / Sink

Flume 的运行单元是 Agent——一个独立的 JVM 进程。每个 Agent 内部由三大组件串成一条流水线:Source(数据从哪来)→ Channel(中间缓冲)→ Sink(数据到哪去)

Flume Agent 内部结构与事务示意图。上方面板"单 Agent 内部:一条 Event 从外部源经 Channel 缓冲到目的地":从左到右五个方块依次为——外部数据源(Nginx / App 日志)、Source(收数据、封装成 Event)、Channel(缓冲队列,可选 Memory 或 File)、Sink(拉取并发往目的地)、目的地(Kafka / HDFS / 下一 Agent);相邻方块之间的箭头标注为"采集""put 事务""take 事务""输出"。下方两个说明框:绿色框讲 put 事务(Source→Channel)——Source 把一批 Event 放入 Channel,全部落定才提交,中途失败则回滚,这批 Event 不会进 Channel;蓝色框讲 take 事务(Channel→Sink)——Sink 取一批 Event 发往目的地,确认成功才从 Channel 删除,失败回滚、Event 留在 Channel 重发,因此不丢但可能重复。

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 ChannelEvent 存内存队列,快、吞吐高;但进程崩溃 / 断电 → 队列数据全丢能容忍少量丢失的场景(如监控指标)
File ChannelEvent 落本地磁盘(含 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 两端各有一个事务(见上图):

  1. put 事务(Source → Channel):Source 把一批 Event 批量放进 Channel,全部成功才 commit;任何一条失败就 rollback,这批 Event 不进 Channel(Source 会重新采集)。
  2. 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 做轻量加工:

3.2 Channel Selector:一份数据分流到多个 Channel

当一个 Source 要把数据发到多个 Channel(进而多个 Sink)时,由 Channel Selector 决定怎么分:

3.3 Sink Processor / Sink Group:多 Sink 的容错与均衡

把多个 Sink 组成一个 Sink Group,由 Sink Processor 决定行为:

这三个组件的组合能力,是 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 → 双消费架构。

广告日志采集分层拓扑图。左侧采集层是三台广告服务器(A/B/C),每台各跑一个 Flume Agent 用 TailDir 采集本地日志;三条箭头标注"Avro RPC 汇聚",扇入到中间的紫色"汇聚 Agent(Avro Source + File Channel)";汇聚 Agent 通过标注"Kafka Sink"的箭头写入黄色的 Kafka(topic 名 ad_raw_log,分区 + 3 副本);Kafka 再通过标注"多消费组"的两条箭头扇出到右侧两个下游——粉色的 Flink(实时计费 / 归因 / 反作弊)与绿色的 HDFS / Hive(离线数仓 T+1)。底部黄色说明框:为什么要汇聚层——把成百上千台机器直连 Kafka 收敛成少数几条连接,减小 Kafka 连接压力,并在汇聚点统一做压缩 / 限流 / 加密;File Channel + acks=-1 是广告计费日志"不丢"的地基。

广告日志采集的经典分层:采集层每机一个 TailDir Agent,经 Avro 汇聚到少数汇聚 Agent,再由 Kafka Sink 统一写入 Kafka,下游 Flink(实时)与 HDFS(离线)各起一个消费组。汇聚层把海量直连收敛成少数连接,也是统一压缩/限流/加密的地方。

这条链路是两个动作的叠加:采集层多 Agent **扇入(fan-in)汇聚层收敛连接,Kafka 再扇出(fan-out)**给 Flink(实时)与 HDFS(离线)两个消费组。

为什么要中间加一层汇聚,而不是每台机器直连 Kafka?

这条链路”不丢”的三个关键点

  1. 采集层用 File Channel——Agent 崩溃重启不丢未发送的 Event;TailDir positionFile 保证不重读不漏读。
  2. Kafka Sink acks=-1 + topic 3 副本 + min.insync.replicas=2——写入侧不丢(原理见 Kafka 可靠性与 Exactly-Once)。
  3. 下游 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 侧 Flume 客户端流程图。一条从左到右的流水线,五个方块依次为:FlumeEventProducer(封 Event + msgType)、LinkedBlockingQueue(内存解耦缓冲)、FlumeEventConsumer(单线程 · 分类攒批)、FlumeRpcClient(loadbalance RR)、Avro Source :4141(multiplexing 分流);相邻方块间的箭头标注为"offer 入队""drainTo(10)""批量""Avro RPC"。底部绿色说明框:msgType 由 Producer 拼成「类型 + 版本」(如 event1_0 / bill1_0),正是 Agent 端 selector.mapping 的路由 key——客户端与服务端靠同一个 header 串起来;队列满则丢弃并告警(保竞价延迟),是客户端侧的一段丢数窗口。

Bidder 侧客户端四段式:业务线程只管把日志 offer 进内存队列(零阻塞),单线程消费者按 msgType 分类攒批,再由 FlumeRpcClient 负载均衡地 Avro RPC 打到汇聚 Agent 的 Source。

FlumeEventProducer(生产者):业务侧调各类 sendXxxMessage,把日志封成 Flume Eventbody = UTF-8 日志正文,headersmsgType / msgVersion / timestamp),非阻塞地 offer 进一个内存队列,超时 20ms 入队失败就记告警。关键在 msgType 的拼法——类型 + 版本,如 event + 1_0event1_0bill1_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)**在这套系统里的真身:

上游 Bidder 横向扩容的扇入拓扑图。左侧纵向排列三个紫色方块"Bidder 实例 1 / 2 / 3 …(几十/上百个),内嵌 FlumeRpcClient",各自通过标注"Avro RPC 扇入(round-robin)"的箭头收敛(多对一)到中间一个绿色方块"少数汇聚 Agent,Avro Source :4141,16 线程";该汇聚 Agent 再通过标注"multiplexing / Sink"的两条箭头发散到右侧两个终点——黄色的"Kafka → Flink(实时)"与绿色的"本地文件 → OSS(离线)"。底部黄色说明框:每个 Bidder 实例内嵌 FlumeRpcClient,靠 round-robin 在多台汇聚 Agent 间负载均衡;Bidder 横向扩多少实例,就有多少个采集端在推;汇聚 Agent 用一个 Avro Source(:4141) 承接上百个客户端连接,正是扇入"收敛连接数"的价值。

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,各按上面两种模式之一映射

映射其实只有两种模式

一句话:竞价链路的十几种事件(出价、竞得、曝光、点击、转化、模型打分……)全从一个源头进来,按类型分流,再按”单路 / 双路”决定后续命运。

6.3 双写:货币化事件同时进「文件(→OSS)」和「Kafka(→实时)」

注意上面 bill1_0 被映射到两个 Channel——这是 §3.2 的 Replicating 效果(一份 Event 复制两份),一份走离线、一份走实时:

DSP 生产案例:一个 bill(计费)事件的双写与 failover 兜底示意图。左侧是"bill 事件(msgType 路由后)",通过标注"复制(双写)"的两条箭头进入中间两个 memory Channel:memChannelBill 与 memChannelKafkaBill;两个 Channel 各自通过"take 事务"箭头连到对应的 Sink:上路是自研 FileSink(落本地文件),下路是一个 failover 组(KafkaSink 优先级 10 → FileSink 优先级 5);两个 Sink 再通过"输出"箭头连到最右侧的两个终点:上路"本地文件 → OSS,离线归档 / 补数",下路"Kafka → Flink,实时计费 / 归因"。底部红色说明框:取舍——全用 Memory Channel + acks=1 换低延迟,靠双写 + Kafka 挂了 failover 落本地盘补 durability;残留风险——进程被硬 kill(OOMKill / 节点故障、跳过优雅停机)时,内存 Channel 里在途事件仍会丢。

一份 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,而是靠一套组合拳:

这是一个**“延迟 / 吞吐优先”的取舍(配合 acks=1snappy、小 batch、linger.ms=100)。但要清醒认识残留风险**:failover 只在”Kafka 不可用”时触发;如果 进程被硬 kill(OOMKill、节点宕机,跳过优雅停机),Memory Channel 里在途的事件会直接丢。能否接受,取决于这些事件有没有离线文件路兜底、以及对账能否补回。

6.6 从这份配置能学到什么

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. 生产实践与踩坑

9. 选型:Flume 还是 Filebeat / Fluentd / Vector?

Flume 只是日志采集 agent 里的一员。放到今天的技术格局里横向对比:

组件语言资源占用特点适用
FlumeJava较高Hadoop 原生,写 HDFS/Kafka 一把梭,可靠性完整存量大数据平台、写 HDFS/Kafka
FilebeatGo轻量,专注采集文件转发,ELK 栈标配每机部署采集本地日志
LogstashJRuby清洗/转换能力强、插件丰富日志需复杂过滤解析
Fluentd / Fluent BitRuby+C / C中 / 极低CNCF 项目,云原生首选,Fluent Bit 极轻K8s DaemonSet 采集容器日志
VectorRust极低新一代,性能强、可观测性好追求性能/现代化的新平台

怎么选(决策路径)

现状与趋势:Flume 社区已进入维护状态(活跃度低、少有新特性),大量新项目转向 Vector / Fluent Bit。但在已建成的 Hadoop 数仓体系里,Flume 依然广泛存活——存量系统不会轻易替换一个”能稳定跑、和 HDFS/Kafka 无缝”的采集层。看你是”老平台”还是”新平台”,是选 Flume 的分水岭。

参考


views
Share this post on:

Previous Post
Flink 深挖(开篇)· 流处理模型与运行时架构
Next Post
为什么 AdTech 偏爱 Kafka:广告事件中枢的五大典型场景与架构设计