第一章.Flume概述

1.定义

Flume是Cloudera提供的一个高可用的,高可靠的,分布式的海量日志采集,聚合和传输的系统,Flume基于流式架构,灵活简单

$01[Flume_概述_入门] - 图1

2.组成架构

$01[Flume_概述_入门] - 图2

  1. # Agent
  2. - Agent是一个JVM进程.它以事件的形式将数据从源头送至目的
  3. - Agent主要有我个部分组成: Source, Channel,Sink
  4. # Source
  5. - Source是负责接收数据到Flume Agent 的组件,Source组件可以处理各种类型,各种格式的日志数据,包括avro,thrift,exec,jms,spooling directory,netcat,taildir,sequence generator,syslog,http,legacy
  6. # Sink
  7. - Sink不断的轮询Channel中的事件且批量的移除它们,并将这些事件批量写入到存储或索引系统,或者被发送到另一个Flume Agent
  8. - Sink组件目的地包括HDFS,logger,avro,thrift,ipc,file,Hbase,solr,自定义
  9. # Channel
  10. - Channel是位于source与sink之间的缓冲区,因此,Channel允许Source和Sink运作在不同的频率上,Channel是线程安全的,可以同时处理几个source的写入操作和几个Sink的读取操作
  11. - Flume自带两种Channel: Memory Channel 和File Channel
  12. - Memory Channel是内存中的队列,Memory Channel在不需要关心数据丢失的情景下适用,如果需要关心数据丢失,那么Memory Channel就不应该使用,因为程序死亡,机器宕机或者重启就会导致数据丢失
  13. - File Channel将所有事件写到磁盘,因此在程序宕机或关闭的情况下不会丢失数据
  14. # Event
  15. - 传输单元,Flume数据传输的基本单元,以Event的形式经数据从源头送至目的地,Event由Header和Body两部分组成,Header用来存放该Event的一些属性,为K-V结构,Body用来存放该条数据,形式为字节数组
Header Body
K-V byte array

第二章.Flume入门

1. Flume安装

安装部署

  1. 将apache-flume-1.9.0-bin.tar.gz上传到linux的/opt/software目录下
  2. 解压apache-flume-1.9.0-bin.tar.gz到/opt/module/目录下
  1. tar -zxvf /opt/software/apache-flume-1.9.0-bin.tar.gz -C /opt/module/
  1. 修改apache-flume-1.9.0-bin的名称为flume-1.9.0
  1. mv /opt/module/apache-flume-1.9.0-bin /opt/module/flume-1.9.0
  1. 配置环境变量
  1. # 打开配置文件
  2. vim /etc/profile.d/my_env.sh
  3. # 添加如下内容
  4. export FLUME_HOME=/opt/module/flume-1.9.0
  5. export PATH=$PATH:$FLUME_HOME/bin
  6. # source 配置文件
  7. source /etc/profile.d/my_env.sh
  1. 将lib文件夹下的guava-11.0.2.jar删除以兼容Hadoop 3.1.3
  1. rm -rf /opt/module/flume-1.9.0/lib/guava-11.0.2.jar

2.Flume入门案例

案例一:监控端口数据官方案例

需求 : 使用Flume监听一个端口,收集该端口数据,并打印到控制台

  1. 安装netcat工具
  1. sudo yum install -y nc
  1. 在flume-1.9.0文件夹下创建目录object/simpleCase/config

$01[Flume_概述_入门] - 图3

  1. 在config目录下创建配置文件flume_1_netcat_logger.conf,并添加如下内容
  1. # 案例1 端口----->logger
  2. # agent
  3. a1.sources = r1
  4. a1.sinks = k1
  5. a1.channels = c1
  6. # source
  7. a1.sources.r1.type = netcat
  8. a1.sources.r1.bind = hadoop102
  9. a1.sources.r1.port = 6666
  10. # sink
  11. a1.sinks.k1.type = logger
  12. # channel
  13. a1.channels.c1.type = memory
  14. a1.channels.c1.capacity = 1000
  15. a1.channels.c1.transactionCapacity = 100
  16. # bind
  17. a1.sources.r1.channels = c1
  18. a1.sinks.k1.channel = c1
  1. 开启Flume监听端口
  1. flume-ng agent --name a1 --conf /opt/module/flume-1.9.0/conf/ --conf-file flume_1_netcat_logger.conf -Dflume.root.logger=INFO,console
  1. 使用netcat工具向本机的6666端口发送内容
  1. nc hadoop102 6666
  2. hello
  3. jcsune
  1. 在flume 监听页面观察接收数据情况

$01[Flume_概述_入门] - 图4

案例二:实时监控单个追加文件

需求:实时监控Hive日志,并上传到HDFS中

需求分析

$01[Flume_概述_入门] - 图5

实现步骤

  1. 确保Hadoop和java环境配置正确
  2. 创建flume_2_exec_hdfs.conf文件,并添加如下内容
  1. # 案例2 实时监控单个追加文件
  2. # agent
  3. a2.sources=r1
  4. a2.sinks=k1
  5. a2.channels=c1
  6. # source
  7. a2.sources.r1.type = exec
  8. a2.sources.r1.command = tail -F /opt/module/flume-1.9.0/object/simpleCase/data/case2.log
  9. a2.sources.r1.shell = /bin/bash -c
  10. # sink
  11. a2.sinks.k1.type = hdfs
  12. a2.sinks.k1.hdfs.path = hdfs://hadoop102:9820/flume-1.9.0/simpleCase/2/%Y%m%d/%H
  13. #上传文件的前缀
  14. a2.sinks.k1.hdfs.filePrefix = logs-
  15. #是否按照时间滚动文件夹
  16. a2.sinks.k1.hdfs.round = true
  17. #多少时间单位创建一个新的文件夹
  18. a2.sinks.k1.hdfs.roundValue = 1
  19. #重新定义时间单位
  20. a2.sinks.k1.hdfs.roundUnit = hour
  21. #是否使用本地时间戳
  22. a2.sinks.k1.hdfs.useLocalTimeStamp = true
  23. #积攒多少个Event才flush到HDFS一次
  24. a2.sinks.k1.hdfs.batchSize = 100
  25. #设置文件类型,可支持压缩
  26. a2.sinks.k1.hdfs.fileType = DataStream
  27. #多久生成一个新的文件
  28. a2.sinks.k1.hdfs.rollInterval = 60
  29. #设置每个文件的滚动大小
  30. a2.sinks.k1.hdfs.rollSize = 134217700
  31. #文件的滚动与Event数量无关
  32. a2.sinks.k1.hdfs.rollCount = 0
  33. # channel
  34. a2.channels.c1.type = memory
  35. a2.channels.c1.capacity = 1000
  36. a2.channels.c1.transactionCapacity = 100
  37. # bind
  38. a2.sources.r1.channels=c1
  39. a2.sinks.k1.channel=c1
  1. 运行Flume
  1. flume-ng agent --name a2 --conf /opt/module/flume-1.9.0/conf/ --conf-file flume_2_exec_hdfs.conf -Dflume.root.logger=INFO,console
  1. 开启Hadoop和Hive并操作Hive产生日志
  1. mycluster.sh start #开启Hadoop
  2. hiveservice.sh start #开启Hive
  1. 测试(/opt/module/flume-1.9.0/object/simpleCase/data)
  1. echo 1 >> case2.log
  2. echo 2 >> case2.log
  3. echo 3 >> case2.log
  4. echo 4 >> case2.log
  5. echo 5 >> case2.log
  1. 浏览器查看

$01[Flume_概述_入门] - 图6

案例三:实时监控目录下多个新文件

案例需求: 使用Flume监听整个目录的文件,并上传至HDFS

$01[Flume_概述_入门] - 图7

实现步骤:

  1. 创建配置文件flume_3_spoolDir_hdfs.conf,并添加如下内容
  1. # 案例3 实时监控目录下多个文件
  2. # agent
  3. a3.sources=r1
  4. a3.sinks=k1
  5. a3.channels=c1
  6. # source
  7. a3.sources.r1.type = spooldir
  8. a3.sources.r1.spoolDir = /opt/module/flume-1.9.0/object/simpleCase/data/
  9. a3.sources.r1.fileSuffix = .finish
  10. a3.sources.r1.fileHeader = true
  11. #忽略所有以.tmp结尾的文件,不上传
  12. a3.sources.r1.ignorePattern = ([^ ]*\.tmp)
  13. # sink
  14. a3.sinks.k1.type = hdfs
  15. a3.sinks.k1.hdfs.path = hdfs://hadoop102:9820/flume-1.9.0/simpleCase/3/%Y%m%d/%H
  16. #上传文件的前缀
  17. a3.sinks.k1.hdfs.filePrefix = logs-
  18. #是否按照时间滚动文件夹
  19. a3.sinks.k1.hdfs.round = true
  20. #多少时间单位创建一个新的文件夹
  21. a3.sinks.k1.hdfs.roundValue = 1
  22. #重新定义时间单位
  23. a3.sinks.k1.hdfs.roundUnit = hour
  24. #是否使用本地时间戳
  25. a3.sinks.k1.hdfs.useLocalTimeStamp = true
  26. #积攒多少个Event才flush到HDFS一次
  27. a3.sinks.k1.hdfs.batchSize = 100
  28. #设置文件类型,可支持压缩
  29. a3.sinks.k1.hdfs.fileType = DataStream
  30. #多久生成一个新的文件
  31. a3.sinks.k1.hdfs.rollInterval = 60
  32. #设置每个文件的滚动大小
  33. a3.sinks.k1.hdfs.rollSize = 134217702
  34. #文件的滚动与Event数量无关
  35. a3.sinks.k1.hdfs.rollCount = 0
  36. # channel
  37. a3.channels.c1.type = memory
  38. a3.channels.c1.capacity = 1000
  39. a3.channels.c1.transactionCapacity = 100
  40. # bind
  41. a3.sources.r1.channels=c1
  42. a3.sinks.k1.channel=c1
  1. 启动监控文件夹命令
  1. flume-ng agent --name a3 --conf /opt/module/flume-1.9.0/conf/ --conf-file flume_3_spoolDir_hdfs.conf -Dflume.root.logger=INFO,console
  1. 在/opt/module/flume-1.9.0/object/simpleCase下新建test文件夹并添加文件

$01[Flume_概述_入门] - 图8

  1. 浏览器端查看

$01[Flume_概述_入门] - 图9

案例四:实时监控目录下多个追加文件

Exec source适用于监控一个实时追加的文件,不能实现断点续传;Spooldir Source适合用于同步新文件,但不适合对实时追加日志的文件进行监听并同步;而Taildir Source适合用于监听多个实时追加的文件,并且能够实现断点续传。

案例需求:使用Flume监听整个目录的实时追加文件,并上传到HDFS

$01[Flume_概述_入门] - 图10

实现步骤:

  1. 创建配置文件flume_4_tailDir_hdfs.conf ,并添加如下内容
  1. # 案例四 实时监控目录下多个追加文件
  2. # agent
  3. a4.sources=r1
  4. a4.sinks=k1
  5. a4.channels=c1
  6. # source
  7. a4.sources.r1.type = TAILDIR
  8. a4.sources.r1.positionFile = /opt/module/flume-1.9.0/tail_dir.json
  9. a4.sources.r1.filegroups = f1 f2
  10. a4.sources.r1.filegroups.f1 = /opt/module/flume-1.9.0/object/simpleCase/file/.*file.*
  11. a4.sources.r1.filegroups.f2 = /opt/module/flume-1.9.0/object/simpleCase/log/.*log.*
  12. # sink
  13. a4.sinks.k1.type = hdfs
  14. a4.sinks.k1.hdfs.path = hdfs://hadoop102:9820/flume-1.9.0/simpleCase/4/%Y%m%d/%H
  15. #上传文件的前缀
  16. a4.sinks.k1.hdfs.filePrefix = case4-
  17. #是否按照时间滚动文件夹
  18. a4.sinks.k1.hdfs.round = true
  19. #多少时间单位创建一个新的文件夹
  20. a4.sinks.k1.hdfs.roundValue = 1
  21. #重新定义时间单位
  22. a4.sinks.k1.hdfs.roundUnit = hour
  23. #是否使用本地时间戳
  24. a4.sinks.k1.hdfs.useLocalTimeStamp = true
  25. #积攒多少个Event才flush到HDFS一次
  26. a4.sinks.k1.hdfs.batchSize = 100
  27. #设置文件类型,可支持压缩
  28. a4.sinks.k1.hdfs.fileType = DataStream
  29. #多久生成一个新的文件
  30. a4.sinks.k1.hdfs.rollInterval = 60
  31. #设置每个文件的滚动大小
  32. a4.sinks.k1.hdfs.rollSize = 134217700
  33. #文件的滚动与Event数量无关
  34. a4.sinks.k1.hdfs.rollCount = 0
  35. # channel
  36. a4.channels.c1.type = memory
  37. a4.channels.c1.capacity = 1000
  38. a4.channels.c1.transactionCapacity = 100
  39. # bind
  40. a4.sources.r1.channels=c1
  41. a4.sinks.k1.channel=c1
  1. 启动监控文件夹命令
  1. flume-ng agent --name a4 --conf /opt/module/flume-1.9.0/conf/ --conf-file flume_4_tailDir_hdfs.conf -Dflume.root.logger=INFO,console
  1. 分别创建好测试文件

$01[Flume_概述_入门] - 图11

$01[Flume_概述_入门] - 图12

  1. 浏览器查看

$01[Flume_概述_入门] - 图13