第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)
| 概念 | 说明 | 注意事项 |
|---|
| Event | Flume 数据传输的基本单位,包含一个字节数组(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.0、Source code repository: https://git.apache.org/flume.git、Compiler: Oracle Java | 若报错”command not found”,检查 PATH 配置。 |
| 检查 Java | java -version | 确保 Java 版本 ≥ 1.8。 |
| 检查帮助 | flume-ng help | 可查看所有子命令(如 agent、source 等)。 |
2.5 单节点 Agent 启动测试
步骤说明
| 步骤 | 说明 | 注意事项 |
|---|
| 创建配置文件 | vim $FLUME_HOME/conf/test-agent.conf | 配置文件名可自定义,建议以 agent 名命名。 |
| 启动 Agent | flume-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.r1、a1.channels.c1。 | 同一 Agent 内 Channel 名不能重复。 |
| 配置文件格式 | 纯文本 .conf 文件,使用 key = value 格式,支持多行(通过反斜杠 \ 连接)。 | 不支持注释嵌套,# 开头为注释。 |
3.2 定义 Sources、Channels、Sinks
| 组件类型 | 语法格式 | 用途 | 注意事项 |
|---|
| Source | agent-name.sources.source-name.type = source-type
agent-name.sources.source-name.param1 = value1 | 定义数据源类型及参数,如监听端口、读取文件等。 | type 必须为 Flume 支持的 Source 类型(如 netcat、exec)。 |
| Channel | agent-name.channels.channel-name.type = channel-type
agent-name.channels.channel-name.param1 = value1 | 定义缓冲区类型及容量、事务大小等。 | type 可为 memory、file、jdbc 等。 |
| Sink | agent-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 → Channel | agent-name.sources.source-name.channels = channel1 channel2 ... | 指定 Source 写入的 Channel 列表。 | 若使用 Multiplexing Selector,可绑定多个 Channel。 |
| Sink → Channel | agent-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 k1 从 c1 读取。 | 必须确保 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
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | source-name.type = avro | 指定 Source 类型为 Avro。 | a1.sources.r1.type = avro | 必须设置。 |
| bind | source-name.bind = hostname | 绑定监听的主机名或 IP。 | a1.sources.r1.bind = 0.0.0.0 | 使用 0.0.0.0 可监听所有接口。 |
| port | source-name.port = port-number | 指定监听端口号。 | a1.sources.r1.port = 41414 | 端口需未被占用,建议 > 1024。 |
| channels | source-name.channels = channel-list | 指定写入的 Channel。 | a1.sources.r1.channels = c1 | 必须与定义的 Channel 名称一致。 |
| threads | source-name.threads = N | 设置处理请求的线程数。 | a1.sources.r1.threads = 5 | 提高并发处理能力。 |
4.2 Thrift Source
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | source-name.type = thrift | 指定 Source 类型为 Thrift。 | a1.sources.r1.type = thrift | 需 Thrift 服务支持。 |
| bind | source-name.bind = hostname | 监听的主机地址。 | a1.sources.r1.bind = localhost | - |
| port | source-name.port = port-number | 监听端口。 | a1.sources.r1.port = 44444 | - |
| channels | source-name.channels = channel-list | 绑定 Channel。 | a1.sources.r1.channels = c1 | - |
| threads | source-name.threads = N | 服务线程数。 | a1.sources.r1.threads = 3 | 默认为 1。 |
4.3 Exec Source
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | source-name.type = exec | 使用命令执行方式采集数据。 | a1.sources.r1.type = exec | - |
| command | source-name.command = shell-command | 要执行的命令。 | a1.sources.r1.command = tail -F /var/log/app.log | tail -F 可监控文件滚动。 |
| shell | source-name.shell = /bin/sh -c | 执行命令的 shell 环境。 | a1.sources.r1.shell = /bin/sh -c | 必须包含 -c。 |
| restart | source-name.restart = true|false | 命令失败是否重启。 | a1.sources.r1.restart = true | 建议设为 true。 |
| restartThrottle | source-name.restartThrottle = millis | 重启间隔(毫秒)。 | a1.sources.r1.restartThrottle = 5000 | 避免频繁重启。 |
| logStderr | source-name.logStderr = true|false | 是否记录标准错误。 | a1.sources.r1.logStderr = true | 便于调试。 |
| batchSize | source-name.batchSize = N | 每次读取的事件数。 | a1.sources.r1.batchSize = 100 | 默认为 1。 |
4.4 Spooling Directory Source
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | source-name.type = spooldir | 监控目录中的新文件。 | a1.sources.r1.type = spooldir | - |
| spoolDir | source-name.spoolDir = /path/to/dir | 要监控的目录路径。 | a1.sources.r1.spoolDir = /var/log/spool | 目录必须存在且可读。 |
| fileSuffix | source-name.fileSuffix = .suffix | 成功处理后文件的后缀(默认 .COMPLETED)。 | a1.sources.r1.fileSuffix = .DONE | 防止重复读取。 |
| deletePolicy | source-name.deletePolicy = never|immediate | 是否删除已处理文件。 | a1.sources.r1.deletePolicy = immediate | immediate 可节省空间。 |
| batchSize | source-name.batchSize = N | 每批处理的事件数。 | a1.sources.r1.batchSize = 100 | - |
| ignorePattern | source-name.ignorePattern = regex | 忽略匹配的文件名。 | a1.sources.r1.ignorePattern = ^\..* | 忽略隐藏文件。 |
| trackerDir | source-name.trackerDir = /path | 记录已处理文件的元数据目录。 | a1.sources.r1.trackerDir = /var/flume/spooldir | 需有写权限。 |
4.5 NetCat Source
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | source-name.type = netcat | 监听 TCP 端口接收文本行。 | a1.sources.r1.type = netcat | - |
| bind | source-name.bind = hostname | 绑定地址。 | a1.sources.r1.bind = localhost | - |
| port | source-name.port = port-number | 监听端口。 | a1.sources.r1.port = 44444 | - |
| channels | source-name.channels = channel-list | 输出 Channel。 | a1.sources.r1.channels = c1 | - |
| max-line-length | source-name.max-line-length = N | 最大行长度(字节)。 | a1.sources.r1.max-line-length = 5120 | 防止超长行。 |
4.6 Kafka Source
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | source-name.type = org.apache.flume.source.kafka.KafkaSource | 从 Kafka 读取数据。 | a1.sources.r1.type = org.apache.flume.source.kafka.KafkaSource | 必须使用完整类名。 |
| kafka.bootstrap.servers | source-name.kafka.bootstrap.servers = host:port | Kafka 集群地址。 | a1.sources.r1.kafka.bootstrap.servers = kafka1:9092,kafka2:9092 | - |
| kafka.topics | source-name.kafka.topics = topic1,topic2 | 订阅的主题列表。 | a1.sources.r1.kafka.topics = logs | 支持正则:topic.* |
| kafka.consumer.group.id | source-name.kafka.consumer.group.id = group-name | 消费者组 ID。 | a1.sources.r1.kafka.consumer.group.id = flume-group | 确保组内唯一消费。 |
| batchDurationMillis | source-name.batchDurationMillis = N | 每批拉取最大等待时间。 | a1.sources.r1.batchDurationMillis = 1000 | - |
| batchSize | source-name.batchSize = N | 每批拉取的最大记录数。 | a1.sources.r1.batchSize = 1000 | - |
| channels | source-name.channels = channel-list | 输出 Channel。 | a1.sources.r1.channels = c1 | - |
4.7 自定义 Source 简介
| 概念 | 说明 | 注意事项 |
|---|
| 实现接口 | 继承 AbstractSource 并实现 Configurable 和 EventDrivenSource 接口。 | 需重写 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
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | channel-name.type = memory | 指定 Channel 类型为内存通道。 | a1.channels.c1.type = memory | 必须设置。 |
| capacity | channel-name.capacity = N | 最大可存储的 Event 数量。 | a1.channels.c1.capacity = 10000 | 默认为 100,建议根据内存调整。 |
| transactionCapacity | channel-name.transactionCapacity = N | 单次事务支持的最大 Event 数。 | a1.channels.c1.transactionCapacity = 1000 | Source/Sink 的 batch size 不能超过此值。 |
| byteCapacityBufferPercentage | channel-name.byteCapacityBufferPercentage = percentage | 用于字节容量计算的缓冲百分比。 | a1.channels.c1.byteCapacityBufferPercentage = 20 | 与 byteCapacity 结合使用。 |
| byteCapacity | channel-name.byteCapacity = bytes | Channel 可用的最大字节数(基于堆内存)。 | a1.channels.c1.byteCapacity = 800000 | 默认为 JVM 堆的 80%。 |
5.2 File Channel
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | channel-name.type = file | 使用磁盘文件持久化 Event。 | a1.channels.c1.type = file | 保证数据不丢失。 |
| checkpointDir | channel-name.checkpointDir = /path/to/dir | 存储检查点元数据的目录。 | a1.channels.c1.checkpointDir = /var/flume/checkpoint | 必须有读写权限。 |
| dataDirs | channel-name.dataDirs = /path1,/path2 | 存储数据文件的目录列表。 | a1.channels.c1.dataDirs = /var/flume/data | 支持多目录提升 I/O 性能。 |
| capacity | channel-name.capacity = N | 最大存储 Event 数。 | a1.channels.c1.capacity = 1000000 | 远高于 Memory Channel。 |
| transactionCapacity | channel-name.transactionCapacity = N | 单事务最大 Event 数。 | a1.channels.c1.transactionCapacity = 10000 | 影响 Source/Sink 批处理性能。 |
| writeTimeout | channel-name.writeTimeout = seconds | 写操作超时时间(秒)。 | a1.channels.c1.writeTimeout = 20 | 默认为 20 秒。 |
| checkpointInterval | channel-name.checkpointInterval = millis | 检查点写入间隔(毫秒)。 | a1.channels.c1.checkpointInterval = 30000 | 默认 30 秒,影响恢复速度。 |
5.3 JDBC Channel
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | channel-name.type = jdbc | 使用关系型数据库存储 Event。 | a1.channels.c1.type = jdbc | 需引入 JDBC 驱动。 |
| driver | channel-name.driver = jdbc-driver-class | 数据库驱动类名。 | a1.channels.c1.driver = org.apache.derby.jdbc.EmbeddedDriver | 如 MySQL:com.mysql.jdbc.Driver。 |
| url | channel-name.url = jdbc-url | 数据库连接 URL。 | a1.channels.c1.url = jdbc:derby:/var/flume/db;create=true | Derby 常用于测试。 |
| username | channel-name.username = user | 数据库用户名。 | a1.channels.c1.username = flume | - |
| password | channel-name.password = pass | 数据库密码。 | a1.channels.c1.password = secret | - |
| dialect | channel-name.dialect = DIALECT | 数据库方言(可选)。 | a1.channels.c1.dialect = DERBY | 自动检测时可省略。 |
| capacity | channel-name.capacity = N | 最大 Event 数。 | a1.channels.c1.capacity = 100000 | 受数据库存储限制。 |
⚠️ 注意:JDBC Channel 在 Flume 1.7+ 中已不推荐使用,建议使用 File 或 Kafka Channel 替代。
5.4 Kafka Channel
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | channel-name.type = org.apache.flume.channel.kafka.KafkaChannel | 使用 Kafka 作为 Channel。 | a1.channels.c1.type = org.apache.flume.channel.kafka.KafkaChannel | 需 Kafka 依赖包。 |
| kafka.bootstrap.servers | channel-name.kafka.bootstrap.servers = host:port | Kafka 集群地址。 | a1.channels.c1.kafka.bootstrap.servers = kafka1:9092 | - |
| kafka.topic | channel-name.kafka.topic = topic-name | 使用的 Kafka 主题。 | a1.channels.c1.kafka.topic = flume-channel | 所有 Agent 共享此主题。 |
| parseAsFlumeEvent | channel-name.parseAsFlumeEvent = true|false | 是否解析为 Flume Event 格式。 | a1.channels.c1.parseAsFlumeEvent = true | 设为 false 时仅传输 body。 |
| batchSize | channel-name.batchSize = N | 每批处理 Event 数。 | a1.channels.c1.batchSize = 1000 | - |
| kafka.producer.acks | channel-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 都会收到相同数据。 |
| Multiplexing | source-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 值决定路由目标。 |
| 可选 Channel | selector.optional = ch1 ch2 | 标记某些 Channel 为可选,失败不影响事务。 | a1.sources.r1.selector.optional = c2 | 常用于日志备份通道。 |
第6章:Flume Sink 详解
6.1 HDFS Sink
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | sink-name.type = hdfs | 写入数据到 HDFS。 | a1.sinks.k1.type = hdfs | 必须设置。 |
| hdfs.path | sink-name.hdfs.path = hdfs://nn:port/dir | HDFS 目标路径,支持转义符。 | a1.sinks.k1.hdfs.path = /flume/events/%Y-%m-%d | %Y, %m, %d 自动替换。 |
| hdfs.filePrefix | sink-name.hdfs.filePrefix = prefix | 文件前缀。 | a1.sinks.k1.hdfs.filePrefix = log- | 默认为 FlumeData。 |
| hdfs.fileType | sink-name.hdfs.fileType = DataStream|SequenceFile | 文件类型。 | a1.sinks.k1.hdfs.fileType = DataStream | 文本日志用 DataStream。 |
| hdfs.writeFormat | sink-name.hdfs.writeFormat = Text|Writable | 写入格式。 | a1.sinks.k1.hdfs.writeFormat = Text | 与 fileType 配合使用。 |
| hdfs.rollInterval | sink-name.hdfs.rollInterval = seconds | 滚动新文件的时间间隔(0 表示禁用)。 | a1.sinks.k1.hdfs.rollInterval = 3600 | 每小时一个文件。 |
| hdfs.rollSize | sink-name.hdfs.rollSize = bytes | 文件大小达到后滚动(0 禁用)。 | a1.sinks.k1.hdfs.rollSize = 134217728 | 128MB。 |
| hdfs.rollCount | sink-name.hdfs.rollCount = events | 事件数达到后滚动(0 禁用)。 | a1.sinks.k1.hdfs.rollCount = 1000000 | - |
| hdfs.batchSize | sink-name.hdfs.batchSize = N | 每次写入 HDFS 的事件数。 | a1.sinks.k1.hdfs.batchSize = 1000 | 提高吞吐量。 |
| hdfs.useLocalTimeStamp | sink-name.hdfs.useLocalTimeStamp = true|false | 使用本地时间替换转义符。 | a1.sinks.k1.hdfs.useLocalTimeStamp = true | 否则使用 Event 时间戳。 |
6.2 Logger Sink
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | sink-name.type = logger | 将 Event 内容输出到日志(用于测试)。 | a1.sinks.k1.type = logger | 仅用于调试。 |
| maxBytesToLog | sink-name.maxBytesToLog = N | 最多打印 body 的字节数。 | a1.sinks.k1.maxBytesToLog = 16 | 避免日志过大。 |
6.3 Avro Sink
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | sink-name.type = avro | 发送数据到另一个 Agent 的 Avro Source。 | a1.sinks.k1.type = avro | 用于级联架构。 |
| hostname | sink-name.hostname = host | 目标 Agent 的主机名。 | a1.sinks.k1.hostname = agent2-host | 必须可达。 |
| port | sink-name.port = port | 目标 Agent 的 Avro Source 端口。 | a1.sinks.k1.port = 41414 | - |
| batchSize | sink-name.batchSize = N | 每批发送的事件数。 | a1.sinks.k1.batchSize = 100 | - |
| connectTimeout | sink-name.connectTimeout = millis | 连接超时时间。 | a1.sinks.k1.connectTimeout = 20000 | - |
| requestTimeout | sink-name.requestTimeout = millis | 请求超时时间。 | a1.sinks.k1.requestTimeout = 20000 | - |
6.4 Thrift Sink
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | sink-name.type = thrift | 发送到 Thrift Source。 | a1.sinks.k1.type = thrift | 用法类似 Avro Sink。 |
| hostname | sink-name.hostname = host | 目标主机。 | a1.sinks.k1.hostname = localhost | - |
| port | sink-name.port = port | 目标端口。 | a1.sinks.k1.port = 44444 | - |
| batchSize | sink-name.batchSize = N | 批处理大小。 | a1.sinks.k1.batchSize = 100 | - |
6.5 Kafka Sink
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | sink-name.type = org.apache.flume.sink.kafka.KafkaSink | 写入 Kafka 主题。 | a1.sinks.k1.type = org.apache.flume.sink.kafka.KafkaSink | 需 Kafka 依赖。 |
| kafka.bootstrap.servers | sink-name.kafka.bootstrap.servers = host:port | Kafka 集群地址。 | a1.sinks.k1.kafka.bootstrap.servers = kafka1:9092 | - |
| kafka.topic | sink-name.kafka.topic = topic-name | 目标主题。 | a1.sinks.k1.kafka.topic = logs | 支持 EL 表达式:%{header}。 |
| batchSize | sink-name.batchSize = N | 每批发送数。 | a1.sinks.k1.batchSize = 1000 | 提高吞吐。 |
| requiredAcks | sink-name.requiredAcks = 1|0|-1 | Kafka 确认机制。 | a1.sinks.k1.requiredAcks = 1 | -1 表示 ISR 全部确认。 |
| producer.type | sink-name.producer.type = sync|async | 生产者类型。 | a1.sinks.k1.producer.type = sync | 推荐 sync 保证可靠性。 |
6.6 File Roll Sink
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | sink-name.type = file_roll | 将 Event 写入本地文件系统。 | a1.sinks.k1.type = file_roll | 用于备份或测试。 |
| sink.directory | sink-name.sink.directory = /path | 输出目录。 | a1.sinks.k1.sink.directory = /var/flume/events | 必须存在且可写。 |
| sink.rollInterval | sink-name.sink.rollInterval = seconds | 文件滚动间隔(0 禁用)。 | a1.sinks.k1.sink.rollInterval = 86400 | 每天一个文件。 |
| batchSize | sink-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.name | 如 type = com.example.MyKafkaSink。 |
第7章:Flume 拦截器(Interceptors)
7.1 Timestamp Interceptor
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | interceptor-name.type = timestamp | 在 Event header 中添加时间戳(timestamp 字段)。 | a1.sources.r1.interceptors = i1
a1.sources.r1.interceptors.i1.type = timestamp | 常用于 HDFS 按时间分区。 |
| preserveExisting | interceptor-name.preserveExisting = true|false | 若 header 已有 timestamp,是否保留原值。 | a1.sources.r1.interceptors.i1.preserveExisting = false | 设为 true 可防止覆盖。 |
✅ 说明:自动添加 header["timestamp"] = System.currentTimeMillis()。
7.2 Host Interceptor
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | interceptor-name.type = host | 添加主机名或 IP 到 Event header。 | a1.sources.r1.interceptors.i2.type = host | 用于标识数据来源主机。 |
| useIP | interceptor-name.useIP = true|false | 使用 IP 地址(true)或主机名(false)。 | a1.sources.r1.interceptors.i2.useIP = true | 默认为 true。 |
| hostHeader | interceptor-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
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | interceptor-name.type = static | 添加固定的静态值到 header。 | a1.sources.r1.interceptors.i3.type = static | 用于添加环境、应用名等元数据。 |
| key | interceptor-name.key = header-key | header 的键名。 | a1.sources.r1.interceptors.i3.key = app_name | 必须设置。 |
| value | interceptor-name.value = header-value | header 的值。 | a1.sources.r1.interceptors.i3.value = user-service | 必须设置。 |
| preserveExisting | interceptor-name.preserveExisting = true|false | 若 key 已存在,是否保留原值。 | a1.sources.r1.interceptors.i3.preserveExisting = true | 防止覆盖重要字段。 |
7.4 Regex Filtering Interceptor
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | interceptor-name.type = regex_filter | 根据正则表达式过滤或保留 Event。 | a1.sources.r1.interceptors.i4.type = regex_filter | 用于日志清洗。 |
| regex | interceptor-name.regex = pattern | 匹配 body 内容的正则表达式。 | a1.sources.r1.interceptors.i4.regex = ^\d{4}-\d{2} | 必须设置。 |
| excludeEvents | interceptor-name.excludeEvents = true|false | true:匹配则删除;false:仅保留匹配项。 | a1.sources.r1.interceptors.i4.excludeEvents = false | 控制过滤逻辑。 |
| matchAll | interceptor-name.matchAll = true|false | 是否匹配整个 body(而非部分)。 | a1.sources.r1.interceptors.i4.matchAll = false | 默认为 false。 |
⚠️ 注意:设为 excludeEvents=true 且 regex 匹配成功,则该 Event 被丢弃。
7.5 Search and Replace Interceptor
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | interceptor-name.type = search_replace | 在 Event body 中查找并替换文本。 | a1.sources.r1.interceptors.i5.type = search_replace | 用于敏感信息脱敏或格式标准化。 |
| searchPattern | interceptor-name.searchPattern = regex | 要查找的正则表达式。 | a1.sources.r1.interceptors.i5.searchPattern = \d{3}-\d{3}-\d{4} | 匹配电话号码。 |
| replaceString | interceptor-name.replaceString = text | 替换为目标字符串。 | a1.sources.r1.interceptors.i5.replaceString = XXX-XXX-XXXX | 可为空。 |
| charset | interceptor-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>) 方法。 | 可选择继承 AbstractInterceptor 或 Builder 类简化开发。 |
| 打包部署 | 编译为 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
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | selector.type = replicating | 将每个 Event 复制到所有关联的 Channel。 | a1.sources.r1.selector.type = replicating | 默认行为,无需显式配置。 |
| optional | selector.optional = ch2 ch3 | 标记某些 Channel 为”可选”,其失败不影响事务成功。 | a1.sources.r1.selector.optional = c2 | 用于备份通道或监控通道。 |
✅ 说明:所有必需 Channel 必须写入成功,事务才提交。
8.2 Multiplexing Channel Selector
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | selector.type = multiplexing | 根据 Event header 的值将 Event 路由到不同 Channel。 | a1.sources.r1.selector.type = multiplexing | 实现数据分流。 |
| header | selector.header = header-key | 指定用于路由的 header 键名。 | a1.sources.r1.selector.header = region | 必须设置。 |
| mapping.value1 | selector.mapping.value1 = ch1 ch2 | 当 header 值为 value1 时,发送到 ch1 和 ch2。 | a1.sources.r1.selector.mapping.cn = c1
a1.sources.r1.selector.mapping.us = c2 | 支持多个值映射。 |
| default | selector.default = ch-default | 当 header 值无匹配时,发送到默认 Channel。 | a1.sources.r1.selector.default = c3 | 建议设置,避免事件丢失。 |
| optional | selector.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 设置优先级。 |
| maxpenalty | agent.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
| 参数名 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| type | processor.type = load_balance | 在多个 Sink 间负载均衡地分发 Event。 | a1.sinkgroups.g1.processor.type = load_balance | 提高吞吐和可用性。 |
| backoff | processor.backoff = true|false | 是否启用退避机制(失败 Sink 暂停)。 | a1.sinkgroups.g1.processor.backoff = true | 推荐开启。 |
| selector | processor.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 = 134217728 | 128MB。 |
| 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 → ES | Flume 写 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,000 | Source/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 Sink | hdfs.batchSize | 1000 ~ 10000 | 减少 HDFS RPC 调用。 |
| HDFS Sink | hdfs.rollInterval, rollSize | 3600, 128MB | 平衡文件数量与大小。 |
| Kafka Sink | batchSize | 1000 ~ 5000 | 提高 Kafka 写入效率。 |
| Kafka Sink | requiredAcks | -1 | 等待 ISR 全部确认,保证可靠性。 |
| Avro Sink | batchSize | 100 ~ 1000 | 根据网络带宽调整。 |
| Avro Sink | connectTimeout, requestTimeout | 20000 ms | 避免因短暂网络问题失败。 |
| File Roll Sink | sink.rollInterval | 86400 (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-*.jar | Flume 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=134217728 和 batchSize=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 添加时间戳用于分区。