一、先搞懂:什么是Flume的多级级联?

很多刚接触大数据数据集成的人,一听到“多级级联”就发懵,觉得是很高深的架构。其实你可以把它想象成快递中转:比如你在海南三亚买了一箱芒果,芒果先从三亚的果园(数据源头)送到三亚本地的快递站(第一级Flume),再送到海口的区域中转仓(第二级Flume),最后送到你家所在城市的快递站(第三级Flume),再到你手里(最终存储的大数据平台)。

普通的Flume是“直连”的:果园直接把芒果拉到你家,一旦你家所在城市的快递站爆仓(大数据平台压力大),或者路上堵车(网络波动),芒果就会坏(数据丢了)。而多级级联就是把数据传输拆成了多个“中转环节”,每个环节只负责一段传输,压力分散、容错性也强。

1.1 核心角色:Flume的三大组件

要理解级联,得先知道Flume最基础的三个部分,用快递来对应会很清楚:

  • Source:数据的“收集端”,对应快递的“揽件员”,负责从源头(比如网站的访问日志、APP的埋点数据)把数据“拿”过来。
  • Channel:数据的“暂存区”,对应快递的“周转箱”,揽件员收完芒果,先放进周转箱,再交给下一个环节,避免数据直接丢了。
  • Sink:数据的“输出端”,对应快递的“配送员”,负责把周转箱里的芒果送到下一个快递站(或者最终目的地)。

多级级联的本质,就是把多个Flume实例串起来:前一个Flume的Sink,对接后一个Flume的Source,形成一条传输链路。

二、为什么多级级联能提升数据集成效率?

直接说结论:它解决了普通直连Flume的三个大问题,效率自然就上去了。

2.1 解决“单点压力过大”的问题

如果是直连模式,所有数据都直接往最终的大数据平台(比如HDFS、Kafka)送,一旦平台的写入能力跟不上(比如HDFS的NameNode处理不了太多并发请求),数据就会堵在Flume里,甚至丢数据。

多级级联相当于把压力拆成了“小份”:比如第一级Flume负责收集10台服务器的日志,第二级Flume负责把10台服务器的日志汇总,第三级再负责往HDFS写。每一级只处理自己范围内的压力,不会出现“整个链路因为最后一步卡壳全堵死”的情况。

2.2 解决“跨地域传输不稳定”的问题

很多企业的业务服务器在多个城市(比如北京、上海、广州都有),如果所有城市的日志都直接往总部的大数据平台送,跨地域的网络波动会导致数据丢包、传输慢。

多级级联就可以在每个城市设一个中转Flume:本地服务器的日志先传到本地的Flume,本地Flume再汇总后往总部送。本地传输的网络稳定,跨地域的传输压力也小,速度自然快。

2.3 解决“数据预处理分散”的问题

普通直连模式下,每台业务服务器的Flume都要做数据清洗(比如去掉无效的日志行、统一时间格式),会占用业务服务器的资源,还容易出现不同服务器清洗规则不一致的问题。

多级级联可以把预处理的工作集中到中间的Flume:第一级Flume只负责收集原始数据,中间的Flume统一做清洗、过滤、汇总,最后再往最终平台送。既不占用业务服务器资源,还能保证数据的一致性。

三、手把手做一个多级级联的完整示例

为了让大家能直接上手,我们做一个最常用的场景:“北京、上海两个城市的业务服务器日志,先传到本地的中转Flume,再汇总到总部的Flume,最后存到HDFS”。

3.1 示例前的准备

3.1.1 技术栈说明

所有示例统一使用:Flume 1.9.0(最稳定的版本)、Linux(CentOS 7)、HDFS(Hadoop 2.7.3)

3.1.2 环境要求

需要准备4台Linux服务器:

  • 北京业务服务器(简称“北京业务”):IP 192.168.1.10,负责产生日志
  • 北京中转Flume服务器(简称“北京中转”):IP 192.168.1.11,负责接收北京业务的日志
  • 上海业务服务器(简称“上海业务”):IP 192.168.1.20,负责产生日志
  • 总部中转Flume服务器(简称“总部中转”):IP 192.168.1.30,负责接收北京、上海中转的日志,再存到HDFS

3.1.3 基础配置

所有服务器都要先安装Flume,安装步骤很简单:

  1. 下载Flume安装包:
wget https://archive.apache.org/dist/flume/1.9.0/apache-flume-1.9.0-bin.tar.gz
  1. 解压:
tar -zxf apache-flume-1.9.0-bin.tar.gz -C /opt/
  1. 配置环境变量,编辑/etc/profile,添加:
export FLUME_HOME=/opt/apache-flume-1.9.0-bin
export PATH=$PATH:$FLUME_HOME/bin
  1. 生效环境变量:
source /etc/profile

3.2 第一级:北京、上海业务服务器的Flume配置

这一级的Flume只做一件事:收集本地业务产生的日志,传到本地的中转Flume。

3.2.1 北京业务服务器的Flume配置

$FLUME_HOME/conf目录下新建配置文件beijing_business.conf

# 定义Agent的名字:beijing_business
agent.name = beijing_business

# 定义Source:taildir,实时监控日志文件的变化
agent.sources = s1
agent.sources.s1.type = org.apache.flume.source.taildir.TaildirSource
# 监控的日志文件路径,比如业务日志存在/var/log/business/
agent.sources.s1.filegroups = f1
agent.sources.s1.filegroups.f1 = /var/log/business/*.log
# 记录文件读取的位置,避免重启后重复读
agent.sources.s1.positionFile = /opt/apache-flume-1.9.0-bin/taildir_position.json

# 定义Channel:memory,暂存数据,速度快
agent.channels = c1
agent.channels.c1.type = memory
agent.channels.c1.capacity = 1000 # 最多暂存1000条数据
agent.channels.c1.transactionCapacity = 100 # 每次传输100条数据

# 定义Sink:avro,把数据传到北京中转的Flume
agent.sinks = k1
agent.sinks.k1.type = avro
# 北京中转Flume的IP和端口
agent.sinks.k1.hostname = 192.168.1.11
agent.sinks.k1.port = 4141

# 把Source、Channel、Sink关联起来
agent.sources.s1.channels = c1
agent.sinks.k1.channel = c1

配置好后,启动这个Flume:

flume-ng agent --name beijing_business --conf $FLUME_HOME/conf --conf-file $FLUME_HOME/conf/beijing_business.conf -Dflume.root.logger=INFO,console

3.2.2 上海业务服务器的Flume配置

和北京的配置几乎一样,只需要改几个参数,在$FLUME_HOME/conf下新建shanghai_business.conf

# 定义Agent的名字:shanghai_business
agent.name = shanghai_business

# 定义Source:taildir,监控上海业务的日志
agent.sources = s1
agent.sources.s1.type = org.apache.flume.source.taildir.TaildirSource
agent.sources.s1.filegroups = f1
agent.sources.s1.filegroups.f1 = /var/log/business/*.log
agent.sources.s1.positionFile = /opt/apache-flume-1.9.0-bin/taildir_position.json

# 定义Channel:memory
agent.channels = c1
agent.channels.c1.type = memory
agent.channels.c1.capacity = 1000
agent.channels.c1.transactionCapacity = 100

# 定义Sink:avro,把数据传到总部中转的Flume
agent.sinks = k1
agent.sinks.k1.type = avro
# 总部中转Flume的IP和端口
agent.sinks.k1.hostname = 192.168.1.30
agent.sinks.k1.port = 4141

# 关联组件
agent.sources.s1.channels = c1
agent.sinks.k1.channel = c1

启动上海业务的Flume:

flume-ng agent --name shanghai_business --conf $FLUME_HOME/conf --conf-file $FLUME_HOME/conf/shanghai_business.conf -Dflume.root.logger=INFO,console

3.3 第二级:北京中转Flume的配置

北京中转的Flume负责接收北京业务的日志,再传到总部中转。在$FLUME_HOME/conf下新建beijing_transfer.conf

# 定义Agent的名字:beijing_transfer
agent.name = beijing_transfer

# 定义Source:avro,接收北京业务Flume传过来的数据
agent.sources = s1
agent.sources.s1.type = avro
# 监听的端口,和北京业务Sink的端口一致
agent.sources.s1.bind = 192.168.1.11
agent.sources.s1.port = 4141

# 定义Channel:memory
agent.channels = c1
agent.channels.c1.type = memory
agent.channels.c1.capacity = 2000 # 中转可以多暂存一些
agent.channels.c1.transactionCapacity = 200

# 定义Sink:avro,把数据传到总部中转
agent.sinks = k1
agent.sinks.k1.type = avro
agent.sinks.k1.hostname = 192.168.1.30
agent.sinks.k1.port = 4141

# 关联组件
agent.sources.s1.channels = c1
agent.sinks.k1.channel = c1

启动北京中转的Flume:

flume-ng agent --name beijing_transfer --conf $FLUME_HOME/conf --conf-file $FLUME_HOME/conf/beijing_transfer.conf -Dflume.root.logger=INFO,console

3.4 第三级:总部中转Flume的配置

总部中转的Flume负责接收北京、上海传过来的日志,再存到HDFS。在$FLUME_HOME/conf下新建headquarters_transfer.conf

# 定义Agent的名字:headquarters_transfer
agent.name = headquarters_transfer

# 定义Source:avro,监听北京、上海传过来的数据
agent.sources = s1
agent.sources.s1.type = avro
agent.sources.s1.bind = 192.168.1.30
agent.sources.s1.port = 4141

# 定义Channel:file,因为要存到HDFS,数据量可能很大,用文件暂存更安全
agent.channels = c1
agent.channels.c1.type = file
# 文件暂存的路径
agent.channels.c1.dataDirs = /opt/apache-flume-1.9.0-bin/file_channel
agent.channels.c1.capacity = 10000 # 最多暂存10000条数据
agent.channels.c1.transactionCapacity = 1000 # 每次传输1000条数据

# 定义Sink:hdfs,把数据存到HDFS
agent.sinks = k1
agent.sinks.k1.type = hdfs
# HDFS的路径,按天生成文件,避免单个文件太大
agent.sinks.k1.hdfs.path = hdfs://192.168.1.40:9000/logs/%Y-%m-%d
# 生成的文件名前缀
agent.sinks.k1.hdfs.filePrefix = business_log
# 生成的文件名后缀
agent.sinks.k1.hdfs.fileSuffix = .log
# 滚动生成新文件的条件:每10分钟生成一个新文件
agent.sinks.k1.hdfs.rollInterval = 600
# 滚动生成新文件的条件:文件大小达到128M生成新文件
agent.sinks.k1.hdfs.rollSize = 134217728
# 滚动生成新文件的条件:文件里的记录数达到10000生成新文件
agent.sinks.k1.hdfs.rollCount = 10000
# 文件的格式:文本
agent.sinks.k1.hdfs.fileType = DataStream

# 关联组件
agent.sources.s1.channels = c1
agent.sinks.k1.channel = c1

启动总部中转的Flume:

flume-ng agent --name headquarters_transfer --conf $FLUME_HOME/conf --conf-file $FLUME_HOME/conf/headquarters_transfer.conf -Dflume.root.logger=INFO,console

3.5 验证效果

在任意一台业务服务器上,往日志文件里写一条测试数据:

echo "测试日志:2024-05-20 10:00:00 用户123访问了首页" >> /var/log/business/test.log

然后去HDFS上查看,应该能看到对应的文件,并且里面有这条测试日志:

hdfs dfs -cat /logs/2024-05-20/business_log.1234567890.log

如果能看到,说明整个级联链路是通的。

四、多级级联的优缺点分析

4.1 优点

  1. 容错性强:如果中间某一级的Flume出问题,只会影响这一段的传输,不会导致整个链路的数据丢失。比如北京中转的Flume挂了,北京业务的日志会暂存在本地的Channel里,等北京中转恢复后,数据会继续传,不会丢。
  2. 传输效率高:跨地域传输时,本地中转的速度比直接跨地域快很多;而且每一级只处理自己的压力,不会出现单点瓶颈。
  3. 管理方便:可以按地域、按业务拆分Flume,比如北京的业务归北京中转管,上海的归上海中转管,出问题时很容易定位到具体的环节。

4.2 缺点

  1. 架构复杂:比直连模式多了很多Flume实例,需要维护的服务器、配置都变多了,对运维的要求更高。
  2. 延迟增加:数据要经过多个中转环节,比直连模式的延迟要高一些,对实时性要求特别高的场景(比如实时推荐)不太适用。
  3. 资源消耗大:每个Flume实例都需要占用CPU、内存,服务器数量变多,整体的资源消耗也会增加。

五、实际应用场景和注意事项

5.1 实际应用场景

  1. 跨地域业务数据集成:比如全国有多个分公司,每个分公司的业务数据先传到本地的中转Flume,再汇总到总部的大数据平台。
  2. 大规模日志收集:比如互联网公司有上万台业务服务器,直接往总部传压力太大,就用多级级联拆分压力。
  3. 数据预处理:比如需要对数据做清洗、过滤、汇总,把这些工作集中到中间的Flume,避免占用业务服务器的资源。

5.2 注意事项

  1. Channel的选择:如果是中转环节,数据量不大的话可以用memory Channel(速度快);如果是最终往HDFS存的环节,一定要用file Channel(安全,不会因为服务器重启丢数据)。
  2. 端口配置:每个Flume的Source和Sink的端口一定要对应,比如北京业务的Sink端口是4141,北京中转的Source端口也必须是4141,不然连不上。
  3. 网络连通性:所有Flume实例之间的网络必须连通,防火墙要开放对应的端口,比如4141端口要允许跨服务器访问。
  4. 数据一致性:如果需要保证数据不丢,要给Flume配置事务机制,比如Sink的transactionCapacity要和Channel的transactionCapacity匹配,避免传输过程中数据丢失。
  5. 监控和告警:因为架构复杂,一定要对每个Flume实例做监控,比如监控Channel的暂存数据量、传输速度、错误日志,一旦某个环节出问题,能及时告警。

六、文章总结

Flume的多级级联架构,本质上是把数据传输的“长链路”拆成了“短链路”,通过分散压力、容错隔离、集中预处理的方式,提升了大规模、跨地域数据集成的效率和稳定性。

虽然它比直连模式复杂,对运维和资源的要求更高,但对于大多数企业的大数据场景来说,它的优势远大于缺点。只要按照我们的示例一步步配置,注意Channel选择、端口匹配、网络连通性这些细节,就能搭建出一个稳定高效的多级级联数据集成链路。