一, 简介

Apache Flume 是一个分布式、高可靠、高可用的海量日志采集、聚合和传输系统,是大数据生态中经典的数据搬运工具Apache Flume.

参考网址:https://flume.liyifeng.org/?flag=fromDoc# 【李易峰】

核心架构(Agent)

Flume 的最小工作单元是 Agent(Java 进程),由三大核心组件组成:Source → Channel → Sink

  • Source(数据源)负责从外部(日志文件、网络端口、Kafka 等)采集数据,封装为 Event(数据传输最小单元),并发送到 Channel。常用类型:Taildir Source(监控文件)、Netcat Source(网络端口)、Avro Source(接收其他 Agent 数据)。

  • Channel(通道)位于 Source 与 Sink 之间的缓冲区,暂存 Event,保障数据安全。常用类型:

    • Memory Channel:内存缓存,速度快,但宕机数据丢失(适合非关键数据)。
    • File Channel:磁盘持久化,可靠性高(生产环境首选)。
    • Kafka Channel:基于 Kafka,高吞吐、高可靠。
  • Sink(数据出口)从 Channel 读取 Event,将数据传输到目标存储(HDFS、Kafka、HBase、Elasticsearch)或下一个 Agent。常用类型:HDFS SinkKafka SinkHBase Sink

三,Flume 数据源(Source)全景介绍

Flume Source 是 Agent 的数据入口,负责从外部系统接收数据、封装为统一 Event 并写入 Channel。按数据来源与传输方式,可分为网络协议、文件采集、系统命令、消息队列、测试生成五大类,以下是核心类型详解与选型建议。

类型 核心能力 适用场景 关键特性
Avro Source 监听 Avro 端口,接收 RPC 事件 跨 Agent 级联、多层拓扑、Avro 客户端发送 高可靠、支持跨节点、配合 Channel 容错
Taildir Source 监控多文件,按偏移量续传 实时采集滚动日志、多文件并行、需断点续传 不修改原文件、支持文件轮换、记录位置
Spooling Directory Source 监听目录,采集新文件 批量日志入库、文件不可修改、需完成后归档 自动重命名 / 删除已读文件、高可靠
Exec Source 执行系统命令,读取 stdout 实时采集命令输出、快速原型、tail -F 场景 轻量、进程退出则停止采集、低可靠
Netcat Source 监听 TCP 端口,按行解析文本 测试调试、简单文本流、快速验证 按换行切分、仅用于测试、低可靠
Syslog TCP/UDP Source 接收 Syslog 协议日志 系统日志、网络设备日志、标准化日志采集 按 Syslog 格式解析、支持 TCP/UDP
Kafka Source 作为消费者从 Kafka Topic 拉取 与 Kafka 生态集成、高吞吐日志 / 数据流 支持 offset 管理、高可用、高吞吐
Sequence Generator Source 生成自增序列数字 测试压测、模拟数据、无需真实业务数据 可控速率、简单轻量、仅测试用

选用建议:

Flume Source 覆盖网络、文件、系统、消息队列等全场景,核心选型逻辑是匹配数据生成方式:实时滚动日志选 Taildir,批量文件选 Spooldir,跨节点选 Avro。配置时需注意 bind 地址、文件路径、偏移量管理等关键参数,避免踩坑。

四 案例分析

一, Avro+Memory+Logger

(给服务器上的一个端口发送消息,消息经过内存,打印到控制台上)

配置文件 (可以通过上面给的网站,根据相对应的数据源,管道,以及数据输出方式进行配置)

1. 在/opt/installs/flume/myconf 目录下 建 avro_memory_logger.conf 文件

a1.sources = r1
a1.channels = c1
a1.sources.r1.type = avro
a1.sources.r1.channels = c1
a1.sources.r1.bind = hadoop11
a1.sources.r1.port = 4141

a1.channels.c1.type = memory

a1.sinks = k1
a1.sinks.k1.type = logger
a1.sinks.k1.channel = c1

2. 启动flume

先启动flume-ng

flume-ng agent -c ../conf -f avro-memory-log.conf -n a1 -Dflume.root.logger=INFO,console

-c  后面跟上 配置文件的路径


-f  跟上自己编写的conf文件


-n  agent的名字


-Dflume.root.logger=INFO,console   INFO 日志输出级别  Debug,INFO,warn,error 等

接着新开一个窗口向端口中发送数据:

flume-ng avro-client -c /opt/installs/flume/conf/ -H localhost -p 4141 -F /home/hivedata/arr1.txt

给avro发消息,使用avro-client

如果在执行启动命令的窗口看到文件的输出内容说明成功

二, Exec + Memory + HDFS

exec 数据源

1. 在/opt/installs/flume/myconf 目录下 建 exec_memory_hdfs.conf 文件

在李易峰中找到响应的配置文件
a1.sources = r1
a1.channels = c1
-- 指定数据源
a1.sources.r1.type = exec  
-- 要跟踪的文件
a1.sources.r1.command = tail -F /home/hivedata/user1.txt
a1.sources.r1.channels = c1

-- 以内存形式建立管道
a1.channels.c1.type = memory

a1.sinks = k1
a1.sinks.k1.type = hdfs
a1.sinks.k1.channel = c1
--输出到文件夹中,以年月日/时分建立文件夹
a1.sinks.k1.hdfs.path = /flume/events/%y-%m-%d/%H%M

-- 给事件添加事件戳
a1.sources.r1.interceptors = i1
a1.sources.r1.interceptors.i1.type = timestamp

2.启动测试  (这次sink 指定的是hdfs 所以要启动hdfs )

首先启动 HDFS

 start-dfs.sh

然后启动

flume-ng agent -c ../myconf -f  exec-memory-hdfs.conf -n a1 -Dflume.root.logger=INFO,console

三,spool+file+hdfs

1. 在/opt/installs/flume/myconf 目录下 建 spool-file-hdfs.conf 文件

-- spool数据源 它会将监视指定目录中产生的新文件,并在新文件出现时从新文件中解析数据出来。并且给文件打上COMPLETE 标签

a1.channels = ch-1
a1.sources = src-1

a1.sources.src-1.type = spooldir
a1.sources.src-1.channels = ch-1
-- 指定监视的目录
a1.sources.src-1.spoolDir = /home/zhangsan

-- 指定管道形式
a1.channels.ch-1.type = file

a1.sinks = k1
a1.sinks.k1.type = hdfs
a1.sinks.k1.channel = ch-1
-- 将信息输出到hdfs上
a1.sinks.k1.hdfs.path = /flume/spool/
-- 指定文件的输出格式 默认是sequencefile 
-- 添加这两段指定文件以文本形式输出
a1.sinks.k1.hdfs.fileType = DataStream
a1.sinks.k1.hdfs.writeFormat = Text

2,启动

flume-ng agent -c ../myconf -f spool-file-hdfs.conf -n a1 -Dflume.root.logger=INFO,console

四, tailDir + Memory + HDFS [ 非常常用 ]

1. 在/opt/installs/flume/myconf 目录下 建 tailDir-memory-hdfs.conf文件

a1.sources = r1
a1.channels = c1
a1.sources.r1.type = TAILDIR
a1.sources.r1.channels = c1
a1.sources.r1.positionFile = /var/log/flume/taildir_position.json
a1.sources.r1.filegroups = f1 f2
a1.sources.r1.filegroups.f1 = /home/zhangsan/.*txt

a1.channels.c1.type = memory

a1.sinks = k1
a1.sinks.k1.type = hdfs
a1.sinks.k1.channel = c1
a1.sinks.k1.hdfs.path = /flume/tailDir/

a1.sinks.k1.hdfs.fileType = DataStream
a1.sinks.k1.hdfs.writeFormat = Text

2. 启动

flume-ng agent -c ../myconf -f tailDir-memory-hdfs.conf -n a1 -Dflume.root.logger=INFO,console

五,选型与避坑指南

1. 选型优先级

  • 生产实时日志:优先 Taildir Source(支持滚动日志、断点续传、不破坏原文件);
  • 批量文件归档:选 Spooling Directory Source(文件完成后归档,高可靠);
  • 跨节点 / 多级拓扑:必选 Avro Source(标准 RPC 协议,跨节点通信)Apache Flume;
  • 测试 / 快速原型:用 Netcat SourceExec Source(配置简单,即开即用);
  • Kafka 生态集成:选 Kafka Source(高吞吐、高可用)。

2. 避坑要点

  1. Avro Source 必配 bind=0.0.0.0:否则仅允许本机连接,外部客户端无法发送数据(你之前的问题核心原因);
  2. Taildir 与 Spooldir 区别:Taildir 监控已存在的追加文件,Spooldir 监控目录中的新文件,根据日志生成方式选择;
  3. Exec Source 避免长期运行命令:进程崩溃会导致数据中断,生产环境需配合监控自愈;
  4. 配置文件路径规范:绝对路径优先,避免相对路径导致的文件找不到问题。

更多推荐