Article

日志采集 Flume

更新于:2026-07-12

第1章:Flume 概述与核心概念

1.1 什么是 Flume

概念名称说明注意事项
Apache Flume一个分布式的、可靠的、高可用的海量日志采集、聚合与传输系统,专为在系统中移动大量日志数据而设计。Flume 最初由 Cloudera 开发,后捐赠给 Apache 基金会,主要用于日志数据的收集,也可用于其他类型的数据流传输。
设计目标支持在日志系统中自动地、可靠地进行数据采集和传输,具备容错性、可扩展性和可管理性。不适合用于低延迟的实时查询场景,而是面向批量或近实时的数据管道。
架构特点基于流式架构,数据以事件(Event)形式在组件间流动,支持多级串联和复杂拓扑。所有数据传输都通过 Agent 内部的事务机制保障可靠性。

1.2 Flume 的应用场景

应用场景说明注意事项
日志聚合从多个 Web 服务器、应用服务器收集日志并集中存储到 HDFS 或 Kafka。适用于结构化或半结构化日志(如 Nginx、Tomcat 日志)。
实时数据管道将日志流实时传输到 Spark Streaming、Flink 或 Kafka 进行处理。需合理配置批处理大小和超时参数以降低延迟。
数据迁移将传统系统日志迁移至大数据平台(如 Hadoop 生态)。需考虑数据格式转换和编码问题。
多源合并从不同目录、不同主机采集数据,统一写入目标系统。可通过多个 Agent 分级采集,避免单点压力。
安全审计收集系统操作日志用于安全分析和合规审计。建议使用 File Channel 保证数据不丢失。

1.3 Flume 的核心组件(Source、Channel、Sink)

组件说明注意事项
Source数据的源头,负责接收或采集外部数据(如日志文件、网络流),并将数据封装为 Event 写入 Channel。Source 类型决定数据接入方式,如 Exec Source 监听文件,Avro Source 接收网络请求。
Channel数据的缓冲区,位于 Source 和 Sink 之间,用于临时存储 Event,保证数据在传输过程中的可靠性。Memory Channel 性能高但不持久;File Channel 持久但占用磁盘。
Sink数据的出口,从 Channel 读取 Event 并写入外部系统(如 HDFS、Kafka)。Sink 成功处理后才会从 Channel 删除 Event,确保不丢失。

1.4 Flume 的数据流模型(Event、Agent)

概念说明注意事项
EventFlume 数据传输的基本单位,包含一个字节数组(body)和一个可选的 header(Map<String, String>)。body 通常为日志行,header 可携带元数据(如 host、timestamp)。header 可用于路由决策(如 Multiplexing Selector)。
Agent一个 JVM 进程,包含 Source、Channel、Sink 三个组件,构成一个独立的数据传输单元。一个物理节点可运行多个 Agent,但通常一个 Agent 足以完成任务。
数据流Source → Channel → Sink,数据在 Agent 内部通过事务方式传输,确保可靠性。数据不经过磁盘时也可在内存中流转(Memory Channel)。

1.5 Flume 的可靠性与事务机制

机制说明注意事项
Source 到 Channel 事务Source 在写入 Channel 前开启事务,写入成功后提交,失败则回滚,确保 Event 不丢失。即使 Source 接收成功,若 Channel 写入失败,仍会重试。
Sink 从 Channel 事务Sink 从 Channel 读取一批 Event,处理成功后提交事务并删除,否则回滚重试。HDFS Sink 写入失败时会回滚,避免数据不一致。
可靠性保障使用 File Channel 可实现持久化存储,即使 Agent 崩溃也能恢复数据。Memory Channel 在崩溃时会丢失未处理数据,不适用于关键场景。
故障转移支持 Failover Sink Processor,当主 Sink 失败时自动切换到备用 Sink。需配置多个 Sink 和 Processor 实现高可用。

第2章:Flume 安装与环境搭建

2.1 系统环境要求

要求项说明注意事项
操作系统Linux(推荐 CentOS、Ubuntu)、macOS、Windows(有限支持)生产环境建议使用 Linux。
Java 版本JDK 1.8 或更高版本必须安装 JDK,仅 JRE 不支持。
内存至少 2GB RAM,建议 4GB 以上若使用 File Channel 或大批次处理,需更多内存。
磁盘空间至少 10GB 可用空间(File Channel 和日志文件)File Channel 会持久化数据到磁盘,需预留空间。
网络能访问外部系统(如 HDFS、Kafka)需配置正确的主机名和端口。

2.2 Flume 的下载与解压

步骤说明注意事项
下载地址官网:https://flume.apache.org/download.html选择稳定版本(如 1.11.0),下载 apache-flume-x.x.x-bin.tar.gz
下载命令示例wget https://downloads.apache.org/flume/1.11.0/apache-flume-1.11.0-bin.tar.gz使用国内镜像加速下载(如清华、阿里云)。
解压命令tar -zxvf apache-flume-1.11.0-bin.tar.gz -C /opt/建议解压到 /opt/usr/local 目录。
安装路径/opt/apache-flume-1.11.0-bin/(示例)可创建软链接便于管理:ln -s apache-flume-1.11.0-bin flume

2.3 环境变量配置

配置项说明注意事项
编辑配置文件vim ~/.bashrc/etc/profile建议所有用户可访问时配置全局 profile。
添加环境变量export FLUME_HOME=/opt/apache-flume-1.11.0-bin
export PATH=$FLUME_HOME/bin:$PATH
路径需根据实际解压位置修改。
生效配置source ~/.bashrc修改后必须重新加载配置文件。
验证变量echo $FLUME_HOME确保输出正确路径。

2.4 验证安装与版本检查

命令说明注意事项
查看版本flume-ng version正常输出应包含 Flume 版本号和编译信息。
预期输出示例Flume 1.11.0Source code repository: https://git.apache.org/flume.gitCompiler: Oracle Java若报错”command not found”,检查 PATH 配置。
检查 Javajava -version确保 Java 版本 ≥ 1.8。
检查帮助flume-ng help可查看所有子命令(如 agent、source 等)。

2.5 单节点 Agent 启动测试

步骤说明

步骤说明注意事项
创建配置文件vim $FLUME_HOME/conf/test-agent.conf配置文件名可自定义,建议以 agent 名命名。
启动 Agentflume-ng agent --conf $FLUME_HOME/conf --name test-agent --conf-file $FLUME_HOME/conf/test-agent.conf -Dflume.root.logger=INFO,console--conf 指定配置目录,--name 与配置文件中 agent 名一致,-D 参数用于控制日志输出。
测试数据发送`echo “Hello Flume”nc localhost 44444`
预期输出Event received: { headers:{} body: 48 65 6C 6C 6F 20 46 6C 75 6D 65 }表示数据已成功采集并处理。

配置内容示例

test-agent.sources = netcat-source
test-agent.channels = memory-channel
test-agent.sinks = logger-sink

test-agent.sources.netcat-source.type = netcat
test-agent.sources.netcat-source.bind = localhost
test-agent.sources.netcat-source.port = 44444

test-agent.channels.memory-channel.type = memory
test-agent.channels.memory-channel.capacity = 1000
test-agent.channels.memory-channel.transactionCapacity = 100

test-agent.sinks.logger-sink.type = logger

test-agent.sources.netcat-source.channels = memory-channel
test-agent.sinks.logger-sink.channel = memory-channel

说明:agent 名为 test-agent,使用 NetCat Source 接收 TCP 数据,Memory Channel 缓冲,Logger Sink 打印日志。

第3章:Flume 配置文件结构详解

3.1 Agent 配置的基本结构

概念说明注意事项
Agent 名称配置中第一个层级为 Agent 的逻辑名称(如 a1),后续组件均隶属于该 Agent。名称不能包含空格或特殊字符,建议使用字母、数字、下划线。
配置层级结构agent-name.sources = source-list
agent-name.channels = channel-list
agent-name.sinks = sink-list
必须先声明组件列表,再分别配置各组件属性。
组件命名每个 Source、Channel、Sink 在其类型内唯一,如 a1.sources.r1a1.channels.c1同一 Agent 内 Channel 名不能重复。
配置文件格式纯文本 .conf 文件,使用 key = value 格式,支持多行(通过反斜杠 \ 连接)。不支持注释嵌套,# 开头为注释。

3.2 定义 Sources、Channels、Sinks

组件类型语法格式用途注意事项
Sourceagent-name.sources.source-name.type = source-type
agent-name.sources.source-name.param1 = value1
定义数据源类型及参数,如监听端口、读取文件等。type 必须为 Flume 支持的 Source 类型(如 netcat、exec)。
Channelagent-name.channels.channel-name.type = channel-type
agent-name.channels.channel-name.param1 = value1
定义缓冲区类型及容量、事务大小等。type 可为 memory、file、jdbc 等。
Sinkagent-name.sinks.sink-name.type = sink-type
agent-name.sinks.sink-name.param1 = value1
定义数据输出目标,如写入 HDFS、Kafka。type 如 hdfs、logger、avro。

3.3 绑定 Source 与 Channel、Sink 与 Channel

绑定类型语法用途注意事项
Source → Channelagent-name.sources.source-name.channels = channel1 channel2 ...指定 Source 写入的 Channel 列表。若使用 Multiplexing Selector,可绑定多个 Channel。
Sink → Channelagent-name.sinks.sink-name.channel = channel-name指定 Sink 从哪个 Channel 读取数据。一个 Sink 只能绑定一个 Channel。
示例配置a1.sources.r1.channels = c1
a1.sinks.k1.channel = c1
将 Source r1 的数据写入 c1,Sink k1c1 读取。必须确保 Channel 名称拼写一致。

3.4 多 Agent 配置示例

串联架构(Agent Chain)

Agent1 的 Sink 发送到 Agent2 的 Source(通常使用 Avro Sink → Avro Source),需确保网络可达和端口开放。

Agent1 配置:

tier1.sources = netcat-source
tier1.channels = memory-channel
tier1.sinks = avro-sink

tier1.sources.netcat-source.type = netcat
tier1.sources.netcat-source.bind = localhost
tier1.sources.netcat-source.port = 44444

tier1.channels.memory-channel.type = memory

tier1.sinks.avro-sink.type = avro
tier1.sinks.avro-sink.hostname = localhost
tier1.sinks.avro-sink.port = 55555

tier1.sources.netcat-source.channels = memory-channel
tier1.sinks.avro-sink.channel = memory-channel

Agent1 接收数据并通过 Avro 发送到 Agent2。

Agent2 配置:

tier2.sources = avro-source
tier2.channels = file-channel
tier2.sinks = logger-sink

tier2.sources.avro-source.type = avro
tier2.sources.avro-source.bind = localhost
tier2.sources.avro-source.port = 55555

tier2.channels.file-channel.type = file

tier2.sinks.logger-sink.type = logger

tier2.sources.avro-source.channels = file-channel
tier2.sinks.logger-sink.channel = file-channel

Agent2 接收数据并打印日志。

3.5 配置文件的验证与调试

方法说明注意事项
启动时检查启动命令中加入 -n agent-name--conf-file,Flume 会自动校验语法。配置错误会直接报错并退出。
日志查看查看 Flume 日志文件(默认输出到控制台,可通过 -Dflume.root.logger=INFO,LOGFILE 重定向)。日志中会显示组件初始化、连接状态等信息。
调试参数启动时添加 -Dflume.root.logger=DEBUG,console 可输出详细调试信息。生产环境不建议使用 DEBUG 级别。
配置测试工具无内置验证命令,但可通过 flume-ng agent --conf conf --name a1 --conf-file test.conf --dry-run 模拟(部分版本支持)。--dry-run 并非所有版本都支持,建议直接启动测试。

第4章:Flume Source 详解

4.1 Avro Source

参数名语法用途代码示例注意事项
typesource-name.type = avro指定 Source 类型为 Avro。a1.sources.r1.type = avro必须设置。
bindsource-name.bind = hostname绑定监听的主机名或 IP。a1.sources.r1.bind = 0.0.0.0使用 0.0.0.0 可监听所有接口。
portsource-name.port = port-number指定监听端口号。a1.sources.r1.port = 41414端口需未被占用,建议 > 1024。
channelssource-name.channels = channel-list指定写入的 Channel。a1.sources.r1.channels = c1必须与定义的 Channel 名称一致。
threadssource-name.threads = N设置处理请求的线程数。a1.sources.r1.threads = 5提高并发处理能力。

4.2 Thrift Source

参数名语法用途代码示例注意事项
typesource-name.type = thrift指定 Source 类型为 Thrift。a1.sources.r1.type = thrift需 Thrift 服务支持。
bindsource-name.bind = hostname监听的主机地址。a1.sources.r1.bind = localhost-
portsource-name.port = port-number监听端口。a1.sources.r1.port = 44444-
channelssource-name.channels = channel-list绑定 Channel。a1.sources.r1.channels = c1-
threadssource-name.threads = N服务线程数。a1.sources.r1.threads = 3默认为 1。

4.3 Exec Source

参数名语法用途代码示例注意事项
typesource-name.type = exec使用命令执行方式采集数据。a1.sources.r1.type = exec-
commandsource-name.command = shell-command要执行的命令。a1.sources.r1.command = tail -F /var/log/app.logtail -F 可监控文件滚动。
shellsource-name.shell = /bin/sh -c执行命令的 shell 环境。a1.sources.r1.shell = /bin/sh -c必须包含 -c
restartsource-name.restart = true|false命令失败是否重启。a1.sources.r1.restart = true建议设为 true。
restartThrottlesource-name.restartThrottle = millis重启间隔(毫秒)。a1.sources.r1.restartThrottle = 5000避免频繁重启。
logStderrsource-name.logStderr = true|false是否记录标准错误。a1.sources.r1.logStderr = true便于调试。
batchSizesource-name.batchSize = N每次读取的事件数。a1.sources.r1.batchSize = 100默认为 1。

4.4 Spooling Directory Source

参数名语法用途代码示例注意事项
typesource-name.type = spooldir监控目录中的新文件。a1.sources.r1.type = spooldir-
spoolDirsource-name.spoolDir = /path/to/dir要监控的目录路径。a1.sources.r1.spoolDir = /var/log/spool目录必须存在且可读。
fileSuffixsource-name.fileSuffix = .suffix成功处理后文件的后缀(默认 .COMPLETED)。a1.sources.r1.fileSuffix = .DONE防止重复读取。
deletePolicysource-name.deletePolicy = never|immediate是否删除已处理文件。a1.sources.r1.deletePolicy = immediateimmediate 可节省空间。
batchSizesource-name.batchSize = N每批处理的事件数。a1.sources.r1.batchSize = 100-
ignorePatternsource-name.ignorePattern = regex忽略匹配的文件名。a1.sources.r1.ignorePattern = ^\..*忽略隐藏文件。
trackerDirsource-name.trackerDir = /path记录已处理文件的元数据目录。a1.sources.r1.trackerDir = /var/flume/spooldir需有写权限。

4.5 NetCat Source

参数名语法用途代码示例注意事项
typesource-name.type = netcat监听 TCP 端口接收文本行。a1.sources.r1.type = netcat-
bindsource-name.bind = hostname绑定地址。a1.sources.r1.bind = localhost-
portsource-name.port = port-number监听端口。a1.sources.r1.port = 44444-
channelssource-name.channels = channel-list输出 Channel。a1.sources.r1.channels = c1-
max-line-lengthsource-name.max-line-length = N最大行长度(字节)。a1.sources.r1.max-line-length = 5120防止超长行。

4.6 Kafka Source

参数名语法用途代码示例注意事项
typesource-name.type = org.apache.flume.source.kafka.KafkaSource从 Kafka 读取数据。a1.sources.r1.type = org.apache.flume.source.kafka.KafkaSource必须使用完整类名。
kafka.bootstrap.serverssource-name.kafka.bootstrap.servers = host:portKafka 集群地址。a1.sources.r1.kafka.bootstrap.servers = kafka1:9092,kafka2:9092-
kafka.topicssource-name.kafka.topics = topic1,topic2订阅的主题列表。a1.sources.r1.kafka.topics = logs支持正则:topic.*
kafka.consumer.group.idsource-name.kafka.consumer.group.id = group-name消费者组 ID。a1.sources.r1.kafka.consumer.group.id = flume-group确保组内唯一消费。
batchDurationMillissource-name.batchDurationMillis = N每批拉取最大等待时间。a1.sources.r1.batchDurationMillis = 1000-
batchSizesource-name.batchSize = N每批拉取的最大记录数。a1.sources.r1.batchSize = 1000-
channelssource-name.channels = channel-list输出 Channel。a1.sources.r1.channels = c1-

4.7 自定义 Source 简介

概念说明注意事项
实现接口继承 AbstractSource 并实现 ConfigurableEventDrivenSource 接口。需重写 process() 方法处理事件生成。
打包部署编译为 JAR 文件,放入 $FLUME_HOME/plugins.d/custom-source/lib/ 目录。目录结构需符合 Flume 插件规范。
配置使用在配置文件中通过 type = fully.qualified.class.name 引用。type = com.example.MySource
依赖管理确保 JAR 包包含所有依赖或使用 shade 打包。避免类冲突。

第5章:Flume Channel 详解

5.1 Memory Channel

参数名语法用途代码示例注意事项
typechannel-name.type = memory指定 Channel 类型为内存通道。a1.channels.c1.type = memory必须设置。
capacitychannel-name.capacity = N最大可存储的 Event 数量。a1.channels.c1.capacity = 10000默认为 100,建议根据内存调整。
transactionCapacitychannel-name.transactionCapacity = N单次事务支持的最大 Event 数。a1.channels.c1.transactionCapacity = 1000Source/Sink 的 batch size 不能超过此值。
byteCapacityBufferPercentagechannel-name.byteCapacityBufferPercentage = percentage用于字节容量计算的缓冲百分比。a1.channels.c1.byteCapacityBufferPercentage = 20与 byteCapacity 结合使用。
byteCapacitychannel-name.byteCapacity = bytesChannel 可用的最大字节数(基于堆内存)。a1.channels.c1.byteCapacity = 800000默认为 JVM 堆的 80%。

5.2 File Channel

参数名语法用途代码示例注意事项
typechannel-name.type = file使用磁盘文件持久化 Event。a1.channels.c1.type = file保证数据不丢失。
checkpointDirchannel-name.checkpointDir = /path/to/dir存储检查点元数据的目录。a1.channels.c1.checkpointDir = /var/flume/checkpoint必须有读写权限。
dataDirschannel-name.dataDirs = /path1,/path2存储数据文件的目录列表。a1.channels.c1.dataDirs = /var/flume/data支持多目录提升 I/O 性能。
capacitychannel-name.capacity = N最大存储 Event 数。a1.channels.c1.capacity = 1000000远高于 Memory Channel。
transactionCapacitychannel-name.transactionCapacity = N单事务最大 Event 数。a1.channels.c1.transactionCapacity = 10000影响 Source/Sink 批处理性能。
writeTimeoutchannel-name.writeTimeout = seconds写操作超时时间(秒)。a1.channels.c1.writeTimeout = 20默认为 20 秒。
checkpointIntervalchannel-name.checkpointInterval = millis检查点写入间隔(毫秒)。a1.channels.c1.checkpointInterval = 30000默认 30 秒,影响恢复速度。

5.3 JDBC Channel

参数名语法用途代码示例注意事项
typechannel-name.type = jdbc使用关系型数据库存储 Event。a1.channels.c1.type = jdbc需引入 JDBC 驱动。
driverchannel-name.driver = jdbc-driver-class数据库驱动类名。a1.channels.c1.driver = org.apache.derby.jdbc.EmbeddedDriver如 MySQL:com.mysql.jdbc.Driver
urlchannel-name.url = jdbc-url数据库连接 URL。a1.channels.c1.url = jdbc:derby:/var/flume/db;create=trueDerby 常用于测试。
usernamechannel-name.username = user数据库用户名。a1.channels.c1.username = flume-
passwordchannel-name.password = pass数据库密码。a1.channels.c1.password = secret-
dialectchannel-name.dialect = DIALECT数据库方言(可选)。a1.channels.c1.dialect = DERBY自动检测时可省略。
capacitychannel-name.capacity = N最大 Event 数。a1.channels.c1.capacity = 100000受数据库存储限制。

⚠️ 注意:JDBC Channel 在 Flume 1.7+ 中已不推荐使用,建议使用 File 或 Kafka Channel 替代。

5.4 Kafka Channel

参数名语法用途代码示例注意事项
typechannel-name.type = org.apache.flume.channel.kafka.KafkaChannel使用 Kafka 作为 Channel。a1.channels.c1.type = org.apache.flume.channel.kafka.KafkaChannel需 Kafka 依赖包。
kafka.bootstrap.serverschannel-name.kafka.bootstrap.servers = host:portKafka 集群地址。a1.channels.c1.kafka.bootstrap.servers = kafka1:9092-
kafka.topicchannel-name.kafka.topic = topic-name使用的 Kafka 主题。a1.channels.c1.kafka.topic = flume-channel所有 Agent 共享此主题。
parseAsFlumeEventchannel-name.parseAsFlumeEvent = true|false是否解析为 Flume Event 格式。a1.channels.c1.parseAsFlumeEvent = true设为 false 时仅传输 body。
batchSizechannel-name.batchSize = N每批处理 Event 数。a1.channels.c1.batchSize = 1000-
kafka.producer.ackschannel-name.kafka.producer.acks = all生产者确认机制。a1.channels.c1.kafka.producer.acks = all推荐设为 all 保证可靠性。

5.5 Channel Selector 配置(Replicating、Multiplexing)

选择器类型语法用途代码示例注意事项
Replicating(默认)source-name.selector.type = replicating
source-name.selector.optional = ch2
将 Event 复制到所有 Channel。a1.sources.r1.selector.type = replicating
a1.sources.r1.selector.optional = c2
所有 Channel 都会收到相同数据。
Multiplexingsource-name.selector.type = multiplexing
source-name.selector.header = header-name
source-name.selector.mapping.value1 = c1
source-name.selector.default = c2
根据 Event header 路由到不同 Channel。a1.sources.r1.selector.type = multiplexing
a1.sources.r1.selector.header = city
a1.sources.r1.selector.mapping.beijing = c1
a1.sources.r1.selector.mapping.shanghai = c2
a1.sources.r1.selector.default = c3
header 值决定路由目标。
可选 Channelselector.optional = ch1 ch2标记某些 Channel 为可选,失败不影响事务。a1.sources.r1.selector.optional = c2常用于日志备份通道。

第6章:Flume Sink 详解

6.1 HDFS Sink

参数名语法用途代码示例注意事项
typesink-name.type = hdfs写入数据到 HDFS。a1.sinks.k1.type = hdfs必须设置。
hdfs.pathsink-name.hdfs.path = hdfs://nn:port/dirHDFS 目标路径,支持转义符。a1.sinks.k1.hdfs.path = /flume/events/%Y-%m-%d%Y, %m, %d 自动替换。
hdfs.filePrefixsink-name.hdfs.filePrefix = prefix文件前缀。a1.sinks.k1.hdfs.filePrefix = log-默认为 FlumeData。
hdfs.fileTypesink-name.hdfs.fileType = DataStream|SequenceFile文件类型。a1.sinks.k1.hdfs.fileType = DataStream文本日志用 DataStream。
hdfs.writeFormatsink-name.hdfs.writeFormat = Text|Writable写入格式。a1.sinks.k1.hdfs.writeFormat = Text与 fileType 配合使用。
hdfs.rollIntervalsink-name.hdfs.rollInterval = seconds滚动新文件的时间间隔(0 表示禁用)。a1.sinks.k1.hdfs.rollInterval = 3600每小时一个文件。
hdfs.rollSizesink-name.hdfs.rollSize = bytes文件大小达到后滚动(0 禁用)。a1.sinks.k1.hdfs.rollSize = 134217728128MB。
hdfs.rollCountsink-name.hdfs.rollCount = events事件数达到后滚动(0 禁用)。a1.sinks.k1.hdfs.rollCount = 1000000-
hdfs.batchSizesink-name.hdfs.batchSize = N每次写入 HDFS 的事件数。a1.sinks.k1.hdfs.batchSize = 1000提高吞吐量。
hdfs.useLocalTimeStampsink-name.hdfs.useLocalTimeStamp = true|false使用本地时间替换转义符。a1.sinks.k1.hdfs.useLocalTimeStamp = true否则使用 Event 时间戳。

6.2 Logger Sink

参数名语法用途代码示例注意事项
typesink-name.type = logger将 Event 内容输出到日志(用于测试)。a1.sinks.k1.type = logger仅用于调试。
maxBytesToLogsink-name.maxBytesToLog = N最多打印 body 的字节数。a1.sinks.k1.maxBytesToLog = 16避免日志过大。

6.3 Avro Sink

参数名语法用途代码示例注意事项
typesink-name.type = avro发送数据到另一个 Agent 的 Avro Source。a1.sinks.k1.type = avro用于级联架构。
hostnamesink-name.hostname = host目标 Agent 的主机名。a1.sinks.k1.hostname = agent2-host必须可达。
portsink-name.port = port目标 Agent 的 Avro Source 端口。a1.sinks.k1.port = 41414-
batchSizesink-name.batchSize = N每批发送的事件数。a1.sinks.k1.batchSize = 100-
connectTimeoutsink-name.connectTimeout = millis连接超时时间。a1.sinks.k1.connectTimeout = 20000-
requestTimeoutsink-name.requestTimeout = millis请求超时时间。a1.sinks.k1.requestTimeout = 20000-

6.4 Thrift Sink

参数名语法用途代码示例注意事项
typesink-name.type = thrift发送到 Thrift Source。a1.sinks.k1.type = thrift用法类似 Avro Sink。
hostnamesink-name.hostname = host目标主机。a1.sinks.k1.hostname = localhost-
portsink-name.port = port目标端口。a1.sinks.k1.port = 44444-
batchSizesink-name.batchSize = N批处理大小。a1.sinks.k1.batchSize = 100-

6.5 Kafka Sink

参数名语法用途代码示例注意事项
typesink-name.type = org.apache.flume.sink.kafka.KafkaSink写入 Kafka 主题。a1.sinks.k1.type = org.apache.flume.sink.kafka.KafkaSink需 Kafka 依赖。
kafka.bootstrap.serverssink-name.kafka.bootstrap.servers = host:portKafka 集群地址。a1.sinks.k1.kafka.bootstrap.servers = kafka1:9092-
kafka.topicsink-name.kafka.topic = topic-name目标主题。a1.sinks.k1.kafka.topic = logs支持 EL 表达式:%{header}
batchSizesink-name.batchSize = N每批发送数。a1.sinks.k1.batchSize = 1000提高吞吐。
requiredAckssink-name.requiredAcks = 1|0|-1Kafka 确认机制。a1.sinks.k1.requiredAcks = 1-1 表示 ISR 全部确认。
producer.typesink-name.producer.type = sync|async生产者类型。a1.sinks.k1.producer.type = sync推荐 sync 保证可靠性。

6.6 File Roll Sink

参数名语法用途代码示例注意事项
typesink-name.type = file_roll将 Event 写入本地文件系统。a1.sinks.k1.type = file_roll用于备份或测试。
sink.directorysink-name.sink.directory = /path输出目录。a1.sinks.k1.sink.directory = /var/flume/events必须存在且可写。
sink.rollIntervalsink-name.sink.rollInterval = seconds文件滚动间隔(0 禁用)。a1.sinks.k1.sink.rollInterval = 86400每天一个文件。
batchSizesink-name.batchSize = N批处理大小。a1.sinks.k1.batchSize = 100-

6.7 自定义 Sink 简介

概念说明注意事项
实现接口继承 AbstractSink 并实现 Configurable 接口。重写 process() 方法处理事件写入。
事务管理process() 中开启 Channel 事务,读取 Event 并提交/回滚。必须正确处理异常以保证可靠性。
打包部署编译为 JAR,放入 $FLUME_HOME/plugins.d/custom-sink/lib/目录结构:plugins.d/sink-name/lib/*.jar
配置使用sink-name.type = fully.qualified.class.nametype = com.example.MyKafkaSink

第7章:Flume 拦截器(Interceptors)

7.1 Timestamp Interceptor

参数名语法用途代码示例注意事项
typeinterceptor-name.type = timestamp在 Event header 中添加时间戳(timestamp 字段)。a1.sources.r1.interceptors = i1
a1.sources.r1.interceptors.i1.type = timestamp
常用于 HDFS 按时间分区。
preserveExistinginterceptor-name.preserveExisting = true|false若 header 已有 timestamp,是否保留原值。a1.sources.r1.interceptors.i1.preserveExisting = false设为 true 可防止覆盖。

说明:自动添加 header["timestamp"] = System.currentTimeMillis()

7.2 Host Interceptor

参数名语法用途代码示例注意事项
typeinterceptor-name.type = host添加主机名或 IP 到 Event header。a1.sources.r1.interceptors.i2.type = host用于标识数据来源主机。
useIPinterceptor-name.useIP = true|false使用 IP 地址(true)或主机名(false)。a1.sources.r1.interceptors.i2.useIP = true默认为 true。
hostHeaderinterceptor-name.hostHeader = header-name自定义 header 键名。a1.sources.r1.interceptors.i2.hostHeader = src_host默认为 host。

示例输出 header{src_host=192.168.1.100}

7.3 Static Interceptor

参数名语法用途代码示例注意事项
typeinterceptor-name.type = static添加固定的静态值到 header。a1.sources.r1.interceptors.i3.type = static用于添加环境、应用名等元数据。
keyinterceptor-name.key = header-keyheader 的键名。a1.sources.r1.interceptors.i3.key = app_name必须设置。
valueinterceptor-name.value = header-valueheader 的值。a1.sources.r1.interceptors.i3.value = user-service必须设置。
preserveExistinginterceptor-name.preserveExisting = true|false若 key 已存在,是否保留原值。a1.sources.r1.interceptors.i3.preserveExisting = true防止覆盖重要字段。

7.4 Regex Filtering Interceptor

参数名语法用途代码示例注意事项
typeinterceptor-name.type = regex_filter根据正则表达式过滤或保留 Event。a1.sources.r1.interceptors.i4.type = regex_filter用于日志清洗。
regexinterceptor-name.regex = pattern匹配 body 内容的正则表达式。a1.sources.r1.interceptors.i4.regex = ^\d{4}-\d{2}必须设置。
excludeEventsinterceptor-name.excludeEvents = true|falsetrue:匹配则删除;false:仅保留匹配项。a1.sources.r1.interceptors.i4.excludeEvents = false控制过滤逻辑。
matchAllinterceptor-name.matchAll = true|false是否匹配整个 body(而非部分)。a1.sources.r1.interceptors.i4.matchAll = false默认为 false。

⚠️ 注意:设为 excludeEvents=true 且 regex 匹配成功,则该 Event 被丢弃。

7.5 Search and Replace Interceptor

参数名语法用途代码示例注意事项
typeinterceptor-name.type = search_replace在 Event body 中查找并替换文本。a1.sources.r1.interceptors.i5.type = search_replace用于敏感信息脱敏或格式标准化。
searchPatterninterceptor-name.searchPattern = regex要查找的正则表达式。a1.sources.r1.interceptors.i5.searchPattern = \d{3}-\d{3}-\d{4}匹配电话号码。
replaceStringinterceptor-name.replaceString = text替换为目标字符串。a1.sources.r1.interceptors.i5.replaceString = XXX-XXX-XXXX可为空。
charsetinterceptor-name.charset = charset-name文本编码(可选)。a1.sources.r1.interceptors.i5.charset = UTF-8默认为平台默认编码。

示例:将日志中的电话号码脱敏。

7.6 自定义 Interceptor 简介

概念说明注意事项
实现接口实现 org.apache.flume.interceptor.Interceptor 接口,包含 initialize()close()List<Event> intercept(List<Event>) 方法。可选择继承 AbstractInterceptorBuilder 类简化开发。
打包部署编译为 JAR 文件,放入 $FLUME_HOME/plugins.d/interceptor-name/lib/ 目录。目录结构需符合 Flume 插件规范。
配置使用在 Source 中通过 interceptors.i1.type = fully.qualified.class.name 引用。type = com.example.SecurityInterceptor
常见用途数据脱敏、字段提取、数据验证、路由标记等。避免在拦截器中执行耗时操作,影响吞吐。

第8章:Flume 通道选择器与处理器

8.1 Replicating Channel Selector

参数名语法用途代码示例注意事项
typeselector.type = replicating将每个 Event 复制到所有关联的 Channel。a1.sources.r1.selector.type = replicating默认行为,无需显式配置。
optionalselector.optional = ch2 ch3标记某些 Channel 为”可选”,其失败不影响事务成功。a1.sources.r1.selector.optional = c2用于备份通道或监控通道。

说明:所有必需 Channel 必须写入成功,事务才提交。

8.2 Multiplexing Channel Selector

参数名语法用途代码示例注意事项
typeselector.type = multiplexing根据 Event header 的值将 Event 路由到不同 Channel。a1.sources.r1.selector.type = multiplexing实现数据分流。
headerselector.header = header-key指定用于路由的 header 键名。a1.sources.r1.selector.header = region必须设置。
mapping.value1selector.mapping.value1 = ch1 ch2当 header 值为 value1 时,发送到 ch1 和 ch2。a1.sources.r1.selector.mapping.cn = c1
a1.sources.r1.selector.mapping.us = c2
支持多个值映射。
defaultselector.default = ch-default当 header 值无匹配时,发送到默认 Channel。a1.sources.r1.selector.default = c3建议设置,避免事件丢失。
optionalselector.optional = ch2指定可选 Channel。a1.sources.r1.selector.optional = c2失败不导致事务回滚。

示例:根据 region=cn 将日志写入本地 HDFS,region=us 写入 Kafka。

8.3 Failover Sink Processor

参数名语法用途代码示例注意事项
sink processor 类型agent.sinkgroups.g1.processor.type = failover实现 Sink 的故障转移,主 Sink 失败时切换到备用。a1.sinkgroups.g1.processor.type = failover保证高可用。
priority.<sink>agent.sinkgroups.g1.processor.priority.k1 = 5
agent.sinkgroups.g1.processor.priority.k2 = 10
数值越小优先级越高(k1 优先于 k2)。a1.sinkgroups.g1.processor.priority.k1 = 5必须为每个 Sink 设置优先级。
maxpenaltyagent.sinkgroups.g1.processor.maxpenalty = millis最大退避时间(毫秒)。a1.sinkgroups.g1.processor.maxpenalty = 30000默认 30 秒。

配置示例:

a1.sinkgroups = g1
a1.sinkgroups.g1.sinks = k1 k2
a1.sinkgroups.g1.processor.type = failover
a1.sinkgroups.g1.processor.priority.k1 = 5
a1.sinkgroups.g1.processor.priority.k2 = 10

8.4 LoadBalancing Sink Processor

参数名语法用途代码示例注意事项
typeprocessor.type = load_balance在多个 Sink 间负载均衡地分发 Event。a1.sinkgroups.g1.processor.type = load_balance提高吞吐和可用性。
backoffprocessor.backoff = true|false是否启用退避机制(失败 Sink 暂停)。a1.sinkgroups.g1.processor.backoff = true推荐开启。
selectorprocessor.selector = round_robin|random负载均衡策略。a1.sinkgroups.g1.processor.selector = round_robin默认为轮询。

配置示例:

a1.sinkgroups = g1
a1.sinkgroups.g1.sinks = k1 k2 k3
a1.sinkgroups.g1.processor.type = load_balance
a1.sinkgroups.g1.processor.selector = round_robin

⚠️ 注意:所有 Sink 必须配置相同的 Channel。

第9章:Flume 与外部系统集成

9.1 Flume 与 HDFS 集成

配置项说明代码示例注意事项
Sink 类型使用 hdfs 类型 Sink 将数据写入 HDFS。a1.sinks.k1.type = hdfs需 Hadoop 客户端环境。
hdfs.path指定 HDFS 路径,支持时间转义符。a1.sinks.k1.hdfs.path = /flume/logs/%Y/%m/%d%Y-%m-%d 按天分区。
hdfs.fileType文件格式:DataStream(文本)、SequenceFile 等。a1.sinks.k1.hdfs.fileType = DataStream日志推荐用 DataStream。
hdfs.writeFormat写入格式,如 Text。a1.sinks.k1.hdfs.writeFormat = Text与 fileType 配合使用。
hdfs.rollInterval文件滚动时间(秒),0 为禁用。a1.sinks.k1.hdfs.rollInterval = 3600每小时生成一个文件。
hdfs.rollSize按大小滚动(字节),0 为禁用。a1.sinks.k1.hdfs.rollSize = 134217728128MB。
hdfs.rollCount按事件数滚动,0 为禁用。a1.sinks.k1.hdfs.rollCount = 1000000百万条事件。
hdfs.batchSize每次刷写 HDFS 的事件数。a1.sinks.k1.hdfs.batchSize = 1000提高吞吐。
Kerberos 支持启用安全认证(可选)。a1.sinks.k1.hdfs.kerberosPrincipal = flume/_HOST@EXAMPLE.COM
a1.sinks.k1.hdfs.kerberosKeytab = /path/to/keytab
集群开启 Kerberos 时需配置。

建议:结合 TimeRoller 和 SizeRoller 实现高效分区与文件管理。

9.2 Flume 与 Kafka 集成

集成方式说明配置示例注意事项
作为 Source(消费 Kafka)使用 Kafka Source 从 Kafka 读取数据。a1.sources.r1.type = org.apache.flume.source.kafka.KafkaSource
a1.sources.r1.kafka.bootstrap.servers = kafka1:9092
a1.sources.r1.kafka.topics = topic1
a1.sources.r1.kafka.consumer.group.id = flume-group
需 flume-ng-kafka-source 依赖。
作为 Sink(写入 Kafka)使用 Kafka Sink 将数据发送到 Kafka。a1.sinks.k1.type = org.apache.flume.sink.kafka.KafkaSink
a1.sinks.k1.kafka.bootstrap.servers = kafka1:9092
a1.sinks.k1.kafka.topic = flume-out
a1.sinks.k1.batchSize = 1000
推荐设置 requiredAcks = -1
作为 Channel(中转)使用 Kafka Channel 替代 File/Memory Channel。a1.channels.c1.type = org.apache.flume.channel.kafka.KafkaChannel
a1.channels.c1.kafka.bootstrap.servers = kafka1:9092
a1.channels.c1.kafka.topic = flume-channel
实现高可用与解耦。
序列化Kafka 默认使用字节数组。a1.sinks.k1.kafka.producer.key.serializer=org.apache.kafka.common.serialization.StringSerializer通常无需修改。

优势:Kafka 作为 Channel 可实现多 Source/Sink 共享,提升可靠性。

9.3 Flume 与 Spark Streaming 集成

集成方式说明配置/代码示例注意事项
Avro 集成(推荐)Spark Streaming 使用 FlumeUtils.createStream 接收 Avro 数据。val stream = FlumeUtils.createStream(ssc, "flume-agent-host", 41414)Flume Sink 使用 Avro Sink。
Pull 方式(PollableSource)Spark 主动拉取 Flume 数据(较少用)。需自定义 Source 实现 PollableSource 接口。复杂,不推荐。
通过 Kafka 中转Flume → Kafka → Spark Streaming(最灵活)。Flume Sink 写 Kafka,Spark 使用 KafkaUtils.createDirectStream解耦,支持多消费者。

推荐架构:Flume → Kafka → Spark Streaming,实现高吞吐与容错。

9.4 Flume 与 Elasticsearch 集成

方法说明配置/实现方式注意事项
使用第三方 Sink如 elasticsearch-sink(非官方),需引入 JAR 包。a1.sinks.k1.type = org.appendium.flume.sink.elasticsearch.ElasticSearchSink
a1.sinks.k1.hostNames = es1:9300,es2:9300
a1.sinks.k1.indexName = flume-logs
a1.sinks.k1.indexType = _doc
依赖 Transport Client(已弃用)。
通过 Kafka 中转Flume → Kafka → Logstash/Custom App → ESFlume 写 Kafka,Logstash 消费并写入 ES。更稳定,支持复杂转换。
自定义 Sink实现 ElasticsearchSink,使用 REST High Level Client。继承 AbstractSink,使用 RestHighLevelClient 批量写入。需处理版本兼容与异常重试。
批量写入提高吞吐,减少请求次数。设置 batchSize = 1000,使用 BulkProcessor。避免单次请求过大。

⚠️ 注意:官方 Flume 不直接支持 ES,建议通过 Kafka + Logstash 或自定义 Sink 实现。

第10章:Flume 性能调优与监控

10.1 Channel 容量与批处理调优

参数说明推荐值注意事项
Memory Channel: capacity最大存储 Event 数。100,000 ~ 1,000,000过大可能导致 GC 停顿。
Memory Channel: transactionCapacity单事务事件数。10,000Source/Sink batch size 不能超过此值。
File Channel: capacity磁盘最大 Event 数。1,000,000+受磁盘空间限制。
File Channel: dataDirs多磁盘目录提升 I/O。/disk1/data,/disk2/data分散负载。
Kafka Channel: batchSize批处理大小。1000 ~ 5000提高吞吐。
通用建议Channel 容量应大于 Source 吞吐峰值。——避免 Source 阻塞。

调优目标:确保 Channel 不成为瓶颈,Sink 能及时消费。

10.2 Sink 批处理大小与超时设置

Sink 类型关键参数推荐配置说明
HDFS Sinkhdfs.batchSize1000 ~ 10000减少 HDFS RPC 调用。
HDFS Sinkhdfs.rollInterval, rollSize3600, 128MB平衡文件数量与大小。
Kafka SinkbatchSize1000 ~ 5000提高 Kafka 写入效率。
Kafka SinkrequiredAcks-1等待 ISR 全部确认,保证可靠性。
Avro SinkbatchSize100 ~ 1000根据网络带宽调整。
Avro SinkconnectTimeout, requestTimeout20000 ms避免因短暂网络问题失败。
File Roll Sinksink.rollInterval86400 (1天)避免生成过多小文件。

原则:增大批处理大小可提升吞吐,但增加延迟;需根据业务需求权衡。

10.3 使用 Ganglia 监控 Flume

配置项说明代码示例注意事项
启用监控在 flume-env.sh 中设置监控类型。export JAVA_OPTS="$JAVA_OPTS -Dflume.monitoring.type=ganglia"
export JAVA_OPTS="$JAVA_OPTS -Dflume.monitoring.hosts=ganglia-server:8649"
需 gmetric 服务运行。
监控指标包括 Source 接入量、Channel 队列、Sink 写出量、Event 处理延迟等。自动上报可在 Ganglia Web UI 查看。
依赖包确保 metrics-ganglia JAR 在 classpath。flume-ng-node/lib/metrics-ganglia-*.jarFlume 1.7+ 默认包含。
多 Agent 监控每个 Agent 独立上报。——可对比各节点性能。

替代方案:也可使用 prometheus(通过 metrics-json + exporter)或 JMX。

10.4 日志分析与常见问题排查

问题现象可能原因排查方法解决方案
Source 无法接收数据端口占用、权限不足、命令失败(Exec Source)查看日志 tail -f flume.log`netstat -anp
Channel 积压严重Sink 写出慢于 Source 写入检查 Channel channel.capacity 与 eventPutSuccessCount调大 Sink batchSize;优化 Sink 目标系统性能。
Sink 连接失败目标服务宕机、网络不通、认证失败日志中搜索 Connection refused、Timeout检查目标服务状态、防火墙、Kerberos 配置。
数据丢失使用 Memory Channel + Agent 崩溃日志中查看 channel.put.backoff 或 sink.batchSize > transactionCapacity改用 File/Kafka Channel;确保 batchSize ≤ transactionCapacity。
高 CPU / GC批量过小、频繁日志输出、内存不足jstat -gc 查看 GC;top 查看 CPU调大批处理;减少 DEBUG 日志;增加 JVM 堆内存。
重复数据Spooling Directory 文件未正确标记查看 spool 目录是否有 .COMPLETED 文件确认 fileSuffix 设置;避免手动修改文件。

建议:生产环境使用 File Channel 或 Kafka Channel,避免数据丢失。

第11章:Flume 高可用与容错机制

11.1 多 Agent 串联架构

项目说明配置示例注意事项
架构目标实现数据采集的分层处理与解耦,提升扩展性与可靠性。Agent1(采集)→ Agent2(聚合)→ HDFS/Kafka适用于大规模分布式环境。
Agent1(边缘 Agent)部署在业务服务器,负责本地日志采集。a1.sources.r1.type = exec
a1.sources.r1.command = tail -F /var/log/app.log
a1.sinks.k1.type = avro
a1.sinks.k1.hostname = aggregator-host
a1.sinks.k1.port = 41414
使用 Exec Source 或 Spooling Directory。
Agent2(聚合 Agent)集中接收多个 Agent 数据,统一写入目标系统。a2.sources.r1.type = avro
a2.sources.r1.bind = 0.0.0.0
a2.sources.r1.port = 41414
a2.sinks.k1.type = hdfs
a2.sinks.k1.hdfs.path = /flume/logs/%Y/%m/%d
可配置多个 Sink 实现高可用。
优势减少对中心系统的网络压力,支持横向扩展。——网络中断时,边缘 Agent 可通过 File Channel 缓存数据。

典型场景:Web 服务器 → Flume Agent(每台)→ 中心 Flume 集群 → HDFS。

11.2 故障转移(Failover)配置

配置项说明代码示例注意事项
Sink Group 类型设置为 failover,实现主备切换。a1.sinkgroups.g1.processor.type = failover必须配置优先级。
priority.<sink>数值越小,优先级越高。a1.sinkgroups.g1.processor.priority.k1 = 5
a1.sinkgroups.g1.processor.priority.k2 = 10
k1 为主,k2 为备。
maxpenalty最大退避时间(毫秒)。a1.sinkgroups.g1.processor.maxpenalty = 30000默认 30 秒。

完整配置:

a1.sinkgroups = g1
a1.sinkgroups.g1.sinks = k1 k2
a1.sinkgroups.g1.processor.type = failover
a1.sinkgroups.g1.processor.priority.k1 = 5
a1.sinkgroups.g1.processor.priority.k2 = 10

工作流程:主 Sink(k1)失败 → 暂停 → 切换到备 Sink(k2)→ 主恢复后自动切回。

11.3 负载均衡配置

配置项说明代码示例注意事项
processor.type设置为 load_balance。a1.sinkgroups.g1.processor.type = load_balance提升吞吐与可用性。
selector负载策略:round_robin(轮询)或 random(随机)。a1.sinkgroups.g1.processor.selector = round_robin推荐轮询。
backoff是否启用退避机制(失败 Sink 暂停)。a1.sinkgroups.g1.processor.backoff = true避免持续失败。

完整配置:

a1.sinkgroups = g1
a1.sinkgroups.g1.sinks = k1 k2 k3
a1.sinkgroups.g1.processor.type = load_balance
a1.sinkgroups.g1.processor.selector = round_robin
a1.sinkgroups.g1.processor.backoff = true

优势:在多个 HDFS NameNode、Kafka Broker 或日志服务器间分摊压力。

11.4 File Channel 的持久化保障

特性说明配置建议注意事项
持久化机制将 Event 写入磁盘文件(data files)和检查点(checkpoint)。a1.channels.c1.type = file
a1.channels.c1.checkpointDir = /flume/checkpoint
a1.channels.c1.dataDirs = /flume/data1,/flume/data2
数据不依赖内存。
故障恢复Agent 重启后,从 checkpoint 恢复未处理的 Event。——确保目录有读写权限。
高可靠性即使 Agent 崩溃,数据也不会丢失(除非磁盘损坏)。结合 hdfs.rollSize=134217728batchSize=1000避免事务过大。
性能权衡写入速度低于 Memory Channel,但远高于数据丢失风险。使用 SSD 提升 I/O 性能适用于关键业务日志。
容量管理capacity 可设为百万级,受磁盘空间限制。a1.channels.c1.capacity = 1000000定期清理过期数据。

最佳实践:生产环境推荐使用 File Channel 或 Kafka Channel 保障数据不丢失。

第12章:Flume 实战案例

12.1 日志采集到 HDFS

组件配置说明配置内容
Source监控本地日志文件。a1.sources.r1.type = exec
a1.sources.r1.command = tail -F /var/log/app.log
Channel使用 File Channel 持久化。a1.channels.c1.type = file
a1.channels.c1.checkpointDir = /flume/checkpoint
a1.channels.c1.dataDirs = /flume/data
Sink写入 HDFS,按天分区。a1.sinks.k1.type = hdfs
a1.sinks.k1.hdfs.path = /flume/logs/%Y-%m-%d
a1.sinks.k1.hdfs.fileType = DataStream
a1.sinks.k1.hdfs.writeFormat = Text
a1.sinks.k1.hdfs.rollInterval = 3600
a1.sinks.k1.hdfs.rollSize = 134217728
a1.sinks.k1.hdfs.batchSize = 1000
绑定——a1.sources.r1.channels = c1
a1.sinks.k1.channel = c1

结果:日志按天存储在 HDFS,文件大小约 128MB,适合 Hive 分析。

12.2 实时日志过滤与转发到 Kafka

需求:使用拦截器过滤错误日志,并转发到 Kafka。

组件配置说明配置内容
Source读取日志文件。a1.sources.r1.type = exec
a1.sources.r1.command = tail -F /var/log/app.log
Interceptor过滤包含 “ERROR” 的日志。a1.sources.r1.interceptors = i1
a1.sources.r1.interceptors.i1.type = regex_filter
a1.sources.r1.interceptors.i1.regex = .*ERROR.*
a1.sources.r1.interceptors.i1.excludeEvents = true
Channel使用 Kafka Channel 提升可靠性。a1.channels.c1.type = org.apache.flume.channel.kafka.KafkaChannel
a1.channels.c1.kafka.bootstrap.servers = kafka1:9092
a1.channels.c1.kafka.topic = flume-channel
Sink写入 Kafka 主题。a1.sinks.k1.type = org.apache.flume.sink.kafka.KafkaSink
a1.sinks.k1.kafka.bootstrap.servers = kafka1:9092
a1.sinks.k1.kafka.topic = error-logs
绑定——a1.sources.r1.channels = c1
a1.sinks.k1.channel = c1

用途:供 Spark/Flink 实时告警系统消费。

12.3 多目录日志合并采集

需求:采集 /app/logs/service1/app/logs/service2 两个目录的日志。

组件配置说明配置内容
Source1采集 service1 日志。a1.sources.r1.type = spooldir
a1.sources.r1.spoolDir = /app/logs/service1
Source2采集 service2 日志。a1.sources.r2.type = spooldir
a1.sources.r2.spoolDir = /app/logs/service2
Channel共用 File Channel。a1.channels.c1.type = file
a1.channels.c1.checkpointDir = /flume/checkpoint
a1.channels.c1.dataDirs = /flume/data
Sink统一写入 HDFS。a1.sinks.k1.type = hdfs
a1.sinks.k1.hdfs.path = /flume/logs/%Y-%m-%d
绑定两个 Source 都写入 c1,Sink 从 c1 读取。a1.sources.r1.channels = c1
a1.sources.r2.channels = c1
a1.sinks.k1.channel = c1

优势:实现多服务日志集中管理,便于分析。

12.4 使用拦截器添加业务字段

需求:为每条日志添加应用名、环境、数据中心等静态字段。

配置项说明配置内容
Interceptor 配置使用 Static Interceptor 添加固定字段。a1.sources.r1.interceptors = i1 i2 i3
a1.sources.r1.interceptors.i1.type = static
a1.sources.r1.interceptors.i1.key = app_name
a1.sources.r1.interceptors.i1.value = user-service
a1.sources.r1.interceptors.i2.type = static
a1.sources.r1.interceptors.i2.key = env
a1.sources.r1.interceptors.i2.value = production
a1.sources.r1.interceptors.i3.type = host
a1.sources.r1.interceptors.i3.hostHeader = src_host
效果Event Header 包含:{app_name="user-service", env="production", src_host="web01"}
用途在 HDFS 或 Kafka 中可通过这些字段进行分类查询与路由。-

扩展:可结合 Timestamp Interceptor 添加时间戳用于分区。