日志采集Flume

1.日志采集Flume 1.1 flume 安装 1.1.1 集群规划

1.日志采集Flume

1.1 flume 安装

1.1.1 集群规划

服务器 hadoop102

服务器 hadoop103

服务器 hadoop104

Flume(采集日志)

Flume

Flume

1.1.2 安装部署

(1)将 apache-flume-1.9.0-bin.tar.gz 上传到 linux 的/opt/software 目录下

(2)解压 apache-flume-1.9.0-bin.tar.gz 到/opt/module/目录下

 tar -zxvf /opt/software/apache-flume-1.9.0-bin.tar.gz -C /opt/module/

(3)修改 apache-flume-1.9.0-bin 的名称为 flume

 mv /opt/module/apache-flume-1.9.0-bin /opt/module/flume

(4)将 lib 文件夹下的 guava-11.0.2.jar 删除以兼容 Hadoop 3.1.3

 rm /opt/module/flume/lib/guava-11.0.2.jar

注意:删除 guava-11.0.2.jar 的服务器节点,一定要配置 hadoop 环境变量。否则会报如下异常。

(5)将 flume/conf 下的 flume-env.sh.template 文件修改为 flume-env.sh,并配置 flume-env.sh 文件

 mv flume-env.sh.template flume-env.sh
 vim flume-env.sh

添加以下内容:

 export JAVA_HOME=/opt/module/jdk1.8.0_212

1.2 项目经验之 Flume 组件选型

1.2.1 Source

(1)Taildir Source 相比 Exec Source、Spooling Directory Source 的优势

TailDir Source:断点续传、多目录。Flume1.6 以前需要自己自定义 Source 记录每次读取文件位置,实现断点续传。不会丢数据,但是有可能会导致数据重复。——重要

Exec Source:可以实时搜集数据,但是在 Flume 不运行或者 Shell 命令出错的情况下,数据将会丢失。

Spooling Directory Source:监控目录,支持断点续传。

(2)batchSize 大小如何设置?

答:Event 1K 左右时,500-1000 合适(默认为 100)

1.2.2 Channel

采用 Kafka Channel,省去了 Sink,提高了效率。KafkaChannel 数据存储在 Kafka 里面,所以数据是存储在磁盘中。

(1)file channel :存储于磁盘中。可靠性高,效率低

(2)memory channel :存储于内存中。可靠性低,效率高

(3)Kafka channel :存储于磁盘中。可靠性高,效率高

1.3 日志采集 Flume 配置

1.3.1 Flume 的具体配置如下

(1)在/opt/module/flume/conf 目录下创建 file-flume-kafka.conf 文件

 vim file-flume-kafka.conf

在文件配置如下内容

 #定义组件
 a1.sources = r1
 a1.channels = c1
 ​
 #配置source
 a1.sources.r1.type = TAILDIR
 a1.sources.r1.filegroups = f1
 a1.sources.r1.filegroups.f1 = /opt/module/applog/log/app.*
 a1.sources.r1.positionFile = /opt/module/flume/taildir_position.json
 a1.sources.r1.interceptors =  i1
 a1.sources.r1.interceptors.i1.type = com.starry.flume.interceptor.ETLInterceptor$Builder
 ​
 #配置channel
 a1.channels.c1.type = org.apache.flume.channel.kafka.KafkaChannel
 a1.channels.c1.kafka.bootstrap.servers = hadoop102:9092,hadoop103:9092
 a1.channels.c1.kafka.topic = topic_log
 a1.channels.c1.parseAsFlumeEvent = false
 ​
 #拼接组件
 a1.sources.r1.channels = c1

注意:com.duan.flume.interceptor.ETLInterceptor 是自定义的拦截器的全类名。需要根据用户自定义的拦截器做相应修改。

1.4 Flume 拦截器

1)创建 Maven 工程 flume-interceptor

2)创建包名:com.starry.flume.interceptor

3)在 pom.xml 文件中添加如下配置

 <dependencies>
     <dependency>
         <groupId>org.apache.flume</groupId>
         <artifactId>flume-ng-core</artifactId>
         <version>1.9.0</version>
         <scope>provided</scope>
     </dependency>
 ​
     <dependency>
         <groupId>com.alibaba</groupId>
         <artifactId>fastjson</artifactId>
         <version>1.2.62</version>
     </dependency>
 </dependencies>
 ​
 <build>
     <plugins>
         <plugin>
             <artifactId>maven-compiler-plugin</artifactId>
             <version>2.3.2</version>
             <configuration>
                 <source>1.8</source>
                 <target>1.8</target>
             </configuration>
         </plugin>
         <plugin>
             <artifactId>maven-assembly-plugin</artifactId>
             <configuration>
                 <descriptorRefs>
                     <descriptorRef>jar-with-dependencies</descriptorRef>
                 </descriptorRefs>
             </configuration>
             <executions>
                 <execution>
                     <id>make-assembly</id>
                     <phase>package</phase>
                     <goals>
                         <goal>single</goal>
                     </goals>
                 </execution>
             </executions>
         </plugin>
     </plugins>
 </build>

注意:scope 中 provided 的含义是编译时用该 jar 包。打包时时不用。因为集群上已经存在 flume 的 jar 包。只是本地编译时用一下。

4)在 com.starry.flume.interceptor 包下创建 JSONUtils 类

 package com.starry.flume.interceptor;
 ​
 import com.alibaba.fastjson.JSON;
 import com.alibaba.fastjson.JSONException;
 ​
 public class JSONUtils {
 //    验证数据是否是JSON
     public static boolean isJSONValidate(String log){
         try {
             JSON.parse(log);
             return true;
         }catch (JSONException e){
             return false;
         }
     }
 }

5)在 com.starry.flume.interceptor 包下创建 ETLInterceptor 类

 package com.starry.flume.interceptor;
 ​
 import org.apache.flume.Context;
 import org.apache.flume.Event;
 import org.apache.flume.interceptor.Interceptor;
 ​
 import java.nio.charset.StandardCharsets;
 import java.util.Iterator;
 import java.util.List;
 ​
 public class ETLInterceptor implements Interceptor {
     @Override
     public void initialize() {
 ​
     }
 ​
     @Override
     public Event intercept(Event event) {
 //        取数据进行校验
 //        1.获取数据
         byte[] body = event.getBody();
         String log = new String(body, StandardCharsets.UTF_8);
 //        2.校验
         if (JSONUtils.isJSONValidate(log)) {
             return event;
         } else {
             return null;
         }
     }
 ​
     @Override
     public List<Event> intercept(List<Event> list) {
         Iterator<Event> iterator = list.iterator();
 ​
         while (iterator.hasNext()){
             Event next = iterator.next();
             if(intercept(next)==null){
                 iterator.remove();
             }
         }
         return list;
     }
 ​
 ​
     public static class Builder implements Interceptor.Builder{
         @Override
         public Interceptor build() {
             return new ETLInterceptor();
         }
         @Override
         public void configure(Context context) {
 ​
         }
     }
 ​
 ​
     @Override
     public void close() {
 ​
     }
 }

6)maven 打包

7)需要先将打好的包放入到 hadoop102 的/opt/module/flume/lib 文件夹下面。

8)分发 Flume 到 hadoop103、hadoop104

 xsync flume/

9)分别在 hadoop102、hadoop103 上启动 Flume

 在flume路径下执行
  bin/flume-ng agent --name a1 --conf-file conf/file-flume-kafka.conf &

1.5.测试 flume-Kafka 通道

(1)生成日志

 lg.sh

(2)消费 Kafka 数据,观察控制台是否有数据获取到

 bin/kafka-console-consumer.sh \  --bootstrap-server hadoop102:9092 --from-beginning --topic topic_log

说明:如果获取不到数据,先检查 Kafka、Flume、Zookeeper 是否都正确启动。再检查 Flume 的拦截器代码是否正常。

1.6 日志采集 Flume 启动停止脚本

(1)在/home/duan/bin 目录下创建脚本 f1.sh

 cd /home/duan/bin/
 vim f1.sh

添加以下脚本内容

 #! /bin/bash
 ​
 case $1 in
 "start"){
         for i in hadoop102 hadoop103
         do
                 echo " --------启动 $i 采集flume-------"
                 ssh $i "nohup /opt/module/flume/bin/flume-ng agent --conf-file /opt/module/flume/conf/file-flume-kafka.conf --name a1 -Dflume.root.logger=INFO,LOGFILE >/opt/module/flume/log1.txt 2>&1  &"
         done
 };; 
 "stop"){
         for i in hadoop102 hadoop103
         do
                 echo " --------停止 $i 采集flume-------"
                 ssh $i "ps -ef | grep file-flume-kafka | grep -v grep |awk  '{print \$2}' | xargs -n1 kill -9 "
         done
 ​
 };;
 esac
  • 说明 1:nohup,该命令可以在你退出帐户/关闭终端之后继续运行相应的进程。nohup 就是不挂起的意思,不挂断地运行命令。

  • 说明 2:awk 默认分隔符为空格

  • 说明 3:$2 是在“”双引号内部会被解析为脚本的第二个参数,但是这里面想表达的含义是 awk 的第二个值,所以需要将他转义,用$2 表示。

  • 说明 4:xargs 表示取出前面命令运行的结果,作为后面命令的输入参数。

(2)增加脚本执行权限

 chmod 777 f1.sh

(3)fl 集群启动脚本

 f1.sh start

(4)fl 集群停止脚本

 f1.sh stop


Comment