大数据工具——Flume
一、Flume概念
中文官方文档:https://flume.liyifeng.org/
1.Flume介绍
1、是一个分布式、可靠的、高可用的日志数据采集框架
2、具有数据流的体系结构
3、具有可调整的可靠性和容错性
4、是Hadoop生态中的一个组件
2.Flume设计

1、Flume的最小运行单元是Agent,三大组件:Source,channel,Sink
2、Flume在运行Agent时候,会占用JVM
3、Flume组件
- Source:作用是与数据源进行交互,采集数据,封装成Event,传给Channel
- Event:采集的数据,对数据封装的对象,Event的结构是键值对,键值对用Header来表示,内部封装Body
- channel:将Source传输的Event进行缓存,传输到Sink
- Flow:Event的传输抽象对象
- interceptor(拦截器):对Source或Sink 进行数据拦截过滤
- Selector(选择器):用于Source,可以将不同的Event分发到不同的Channel中
- Sink:接收Channel传过来的Event,然后下沉到对应的储存系统中
- Client:客户端,启动的用户
3.数据模型
- 单一数据模型
整个数据流为
Web Server --> Source --> Channel --> Sink --> HDFS

- 多数据流模型
- 多 Agent 串行传输数据流模型

- 多 Agent 汇聚数据流模型

- 单 Agent 多路数据流模型

- Sinkgroups 数据流模型

4.采集方案模板
- 命名
agentName.sources = s1,s2,s3…
agentName.channels = c1,c2,c3…
agentName.sinks = ss1,ss2,ss3… - 配置三大组件的关联关系
agentName.sources.s1.channels=c1
agentName.sinks.ss1.channels=c1
- 其他
source
channel
sink
二、案例
1、Source & Channel & Sink类型
source
# Avro Source (会)
监听某一个Ip开放的端口,收集数据
# Exec Source (会)
使用系统命令(cat 或者 tail) 监控某一个文件
# Spooling Directory Source (会)
可以监控一个文件夹内的所有文件数据
# Kafka Source(后期学习了Kafka掌握)
监控Kafka消息队列数据流,实时收集数据
# Syslog Sources
通过网络协议进行数据监控(TCP协议和UDP协议)
# HTTP Source
监控HTTP访问URL(类似于监控URL协议数据)
Channel
# Memory Channel (会)
内存保存数据,速度快,但是容易丢失数据(服务器宕机)
# JDBC Channel
目前来说还不是很稳定,因为支持的功能较少
# Kafka Channel (会)
使用Kafka可以保证数据不丢失,而且传输速度也很快,但是比不上内存
# File Channel (会)
将我们的数据落地到本地磁盘中,这样也可以保证数据不丢,但是数据会重复,而且速度慢
Sink
# Logger Sink
将数据打印到控制台,主要是做测试使用
# Arvo Sink (会)
将数据写入某个端口中
# HDFS Sink (会)
将数据写入HDFS中
# Hive Sink
将数据写入Hive中
# Kafka Sink (后期会)
将数据写入到kafka消息队列(目前新版本1.7以上,不支持kafka0.10以下版本)
2.Avro + Memory + Logger
# 命名
a1.sources = r1
a1.channels = c1
a1.sinks = s1
# 关联
a1.sources.r1.channels = c1
a1.sinks.s1.channel = c1
# 设置Source类型和属性
a1.sources.r1.type = avro
a1.sources.r1.bind = 192.168.1.100
a1.sources.r1.port = 10086
# 设置Channel类型和属性 内存会丢失数据的
a1.channels.c1.type = memory
# 缓存池大小
a1.channels.c1.capacity = 1000
# 每个事务sink拉取的大小
a1.channels.c1.transactionCapacity = 100
# 设置Sink类型和属性
a1.sinks.s1.type = logger
a1.sinks.s1.maxBytesToLog = 32
3.实时采集spool +file + hdfs
与Exec区别在于,它可以采集一个目录下面的所有文件,而exec只能采集一个具体的文件
与 Exec Source不同,Spooling Directory Source是可靠的,即使Flume重新启动或被kill,也不会丢失数据。同时作为这种可靠性的代价,指定目录中的被收集的文件必须是不可变的、唯一命名的
# 命名
a1.sources = r1
a1.channels = c1
a1.sinks = s1
# 关联
a1.sources.r1.channels = c1
a1.sinks.s1.channel = c1
# 配置Source的类型和属性
a1.sources.r1.type = spooldir
# 配置监听目录
a1.sources.r1.spoolDir = /usr/local/flume-1.8/flumedata/data
a1.sources.r1.fileSuffix = .COMPLETED
a1.sources.r1.deletePolicy = never
a1.sources.r1.basenameHeader = true
a1.sources.r1.basenameHeaderKey = filename
# 设置Channel类型和属性
a1.channels.c1.type = file
# 配置Sink类型和属性
a1.sinks.s1.type = hdfs
a1.sinks.s1.hdfs.path = hdfs://hdp01:9000/flume/%Y%m%d/%H%M
a1.sinks.s1.hdfs.filePrefix = FlumeTest
a1.sinks.s1.hdfs.fileSuffix = .bk
# 下面三个配置参数如果都设置为0,那么表示不执行次参数(失效)
a1.sinks.s1.hdfs.rollInterval = 100
a1.sinks.s1.hdfs.rollSize = 1000
a1.sinks.s1.hdfs.rollCount = 100
# 设置采集文件格式 如果你是存文本文件,就是用DataStream
a1.sinks.s1.hdfs.fileType = DataStream
a1.sinks.s1.hdfs.writeFormat = Text
# 开启本地时间戳获取参数,因为我们的目录上面已经使用转义符号,所以要使用时间戳
a1.sinks.s1.hdfs.useLocalTimeStamp = true
4.Syslog+Memory+Logger
# 命名
a1.sources = r1
a1.channels = c1
a1.sinks = s1
# 关联
a1.sources.r1.channels = c1
a1.sinks.s1.channel = c1
# 设置Source类型和属性
a1.sources.r1.type = syslogtcp
a1.sources.r1.bind = hdp01
a1.sources.r1.port = 10086
# 设置Channel类型和属性 内存会丢失数据的
a1.channels.c1.type = memory
# 缓存池大小
a1.channels.c1.capacity = 1000
# 每个事务sink拉取的大小
a1.channels.c1.transactionCapacity = 100
# 设置Sink类型和属性
a1.sinks.s1.type = logger
a1.sinks.s1.maxBytesToLog = 32
5.taildir + memory +HDFS
在1.6以后的版本才发行的新特性,以前是没有的
Taildir Source是可靠的,即使发生文件轮换也不会丢失数据。它会定期地以JSON格式在一个专门用于定位的文件上记录每个文件的最后读取位置。如果Flume由于某种原因停止或挂掉,它可以从文件的标记位置重新开始读取。
spooldir和taildir的相同点和不同点
相同点:
1.都是可靠源
2.监听的都是目录里的文件
3.目录都要提前存在
4.文件名不能重复
不同点:
1.spooldir读完后,会修改文件的名称,添加后缀
2.spooldir采集的是新文件
3.taildir监听的文件不会更名,可以一直监听文件尾部的新数据
# 命名
a1.sources = r1
a1.channels = c1
a1.sinks = s1
# 关联
a1.sources.r1.channels = c1
a1.sinks.s1.channel = c1
# 配置source类型和属性
a1.sources.r1.type = TAILDIR
a1.sources.r1.filegroups = g1 g2
a1.sources.r1.filegroups.g1 = /usr/local/flume-1.8/flumedata/test1/.*.txt
a1.sources.r1.filegroups.g2 = /usr/local/flume-1.8/flumedata/test2/.*.csv
# 元数据保存位置
a1.sources.r1.positionFile = /usr/local/flume-1.8/fluem-log/taildir_position.json
# 配置channel类型属性
a1.channels.c1.type = memory
# 缓存池大小
a1.channels.c1.capacity = 1000
# 每个事务sink拉取的大小
a1.channels.c1.transactionCapacity = 100
# 配置Sink类型和属性
a1.sinks.s1.type = hdfs
a1.sinks.s1.hdfs.path = hdfs://hdp01:9000/flume/%Y%m%d/%H%M
a1.sinks.s1.hdfs.filePrefix = FlumeTest
a1.sinks.s1.hdfs.fileSuffix = .bk
# 下面三个配置参数如果都设置为0,那么表示不执行次参数(失效)
a1.sinks.s1.hdfs.rollInterval = 30
a1.sinks.s1.hdfs.rollSize = 1000
a1.sinks.s1.hdfs.rollCount = 100
# 设置采集文件格式 如果你是存文本文件,就是用DataStream
a1.sinks.s1.hdfs.fileType = DataStream
a1.sinks.s1.hdfs.writeFormat = Text
# 开启本地时间戳获取参数,因为我们的目录上面已经使用转义符号,所以要使用时间戳
a1.sinks.s1.hdfs.useLocalTimeStamp = true
三、自动容灾及负载均衡
负载均衡就方式是把channel里面的Event按照配置的负载机制(比如轮询)分别发送到sink各自对应的目的地;故障转移就是这N个sink同一时间只有一个在工作,其余的作为备用,工作的sink挂掉之后备用的sink顶上。
当我们的Channel传输数据到Sink的时候,如果Sink挂掉了,此时可以配置故障转移,启动多个Sink,其中有一个Sink是使用状态,其他Sink等待使用,如果使用的Sink挂掉了,那么等待的Sink就会顶替上,但是要根据优先级,考虑哪个Sink顶替上。
文件配置
# 命名
a1.sources = r1
a1.channels = c1
a1.sinks = s1 s2
# 关联
a1.sources.r1.channels = c1
a1.sinks.s1.channel = c1
a1.sinks.s2.channel = c1
# Source
a1.sources.r1.type = syslogtcp
a1.sources.r1.bind = hdp01
a1.sources.r1.port = 10086
# Channel
a1.channels.c1.type = memory
a1.channels.c1.capacity = 1000
a1.channels.c1.transactionCapacity = 100
# Sink
a1.sinks.s1.type = avro
a1.sinks.s1.hostname = hdp02
a1.sinks.s1.port = 10087
a1.sinks.s2.type = avro
a1.sinks.s2.hostname = hdp03
a1.sinks.s2.port = 10088
# 设置Sink组
a1.sinkgroups = g1
a1.sinkgroups.g1.sinks = s1 s2
a1.sinkgroups.g1.processor.type = failover
a1.sinkgroups.g1.processor.priority.s1 = 50
a1.sinkgroups.g1.processor.priority.s2 = 10
a1.sinkgroups.g1.processor.maxpenalty = 10000
下游Agent,在hdp02上实现
# 命名
a1.sources = r1
a1.channels = c1
a1.sinks = s1
# 关联
a1.sources.r1.channels = c1
a1.sinks.s1.channel = c1
# 设置Source类型和属性
a1.sources.r1.type = avro
a1.sources.r1.bind = hdp02
a1.sources.r1.port = 10087
# 设置Channel类型和属性 内存会丢失数据的
a1.channels.c1.type = memory
# 缓存池大小
a1.channels.c1.capacity = 1000
# 每个事务sink拉取的大小
a1.channels.c1.transactionCapacity = 100
# 设置Sink类型和属性
a1.sinks.s1.type = logger
a1.sinks.s1.maxBytesToLog = 32
下游Agent,在hdp03上实现
# 命名
a1.sources = r1
a1.channels = c1
a1.sinks = s1
# 关联
a1.sources.r1.channels = c1
a1.sinks.s1.channel = c1
# 设置Source类型和属性
a1.sources.r1.type = avro
a1.sources.r1.bind = hdp03
a1.sources.r1.port = 10088
# 设置Channel类型和属性 内存会丢失数据的
a1.channels.c1.type = memory
# 缓存池大小
a1.channels.c1.capacity = 1000
# 每个事务sink拉取的大小
a1.channels.c1.transactionCapacity = 100
# 设置Sink类型和属性
a1.sinks.s1.type = logger
a1.sinks.s1.maxBytesToLog = 32
更多推荐
所有评论(0)