欢迎光临
我们一直在努力

12-学习笔记尚硅谷数仓搭建-将Kafka中的日志数据同步到Hadoop集群的HDFS分布式文件系统的flume配置

目录

 

一、Kafka数据同步到HDFS的flume配置

二、解决日志数据零点漂移问题

三、编写拦截器

四、编写flume日志服务脚本

五、使用拦截器并将数据传输到HDFS

六、在HDFS上查看数据

备注:没有特别说明都在atguigu用户执行命令

一、Kafka数据同步到HDFS的flume配置

进入job目录创建flume配置(hadoop102执行下面命令)

cd /opt/module/flume/job
vim kafka_to_hdfs_log.conf

添加下面的内容:

#定义组件
a1.sources=r1
a1.channels=c1
a1.sinks=k1

#配置source1
a1.sources.r1.type = org.apache.flume.source.kafka.KafkaSource
a1.sources.r1.batchSize = 5000
a1.sources.r1.batchDurationMillis = 2000
a1.sources.r1.kafka.bootstrap.servers = hadoop102:9092,hadoop103:9092,hadoop104:9092
a1.sources.r1.kafka.topics=topic_log
a1.sources.r1.interceptors = i1
a1.sources.r1.interceptors.i1.type = com.atguigu.gmall.flume.interceptor.TimestampInterceptor$Builder

#配置channel
a1.channels.c1.type = file
a1.channels.c1.checkpointDir = /opt/module/flume/checkpoint/behavior1
a1.channels.c1.dataDirs = /opt/module/flume/data/behavior1
a1.channels.c1.maxFileSize = 2146435071
a1.channels.c1.capacity = 1000000
a1.channels.c1.keep-alive = 6

#配置sink
a1.sinks.k1.type = hdfs
a1.sinks.k1.hdfs.path = /origin_data/gmall/log/topic_log/%Y-%m-%d
a1.sinks.k1.hdfs.filePrefix = log
a1.sinks.k1.hdfs.round = false

a1.sinks.k1.hdfs.rollInterval = 10
a1.sinks.k1.hdfs.rollSize = 134217728
a1.sinks.k1.hdfs.rollCount = 0

#控制输出文件类型
a1.sinks.k1.hdfs.fileType = CompressedStream
a1.sinks.k1.hdfs.codeC = gzip

#组装
a1.sources.r1.channels = c1
a1.sinks.k1.channel = c1

二、解决日志数据零点漂移问题

1.出现问题的原因

因为我们是通过Kafka的event的header中的时间戳作为日志数据分区的依据,所以导致一个在当天23:59:59的日志文件在经过传输后到达Kafka的时间变为第二天的0:0:1,这样就错误的将当天的日志数据划分到了第二天

2.解决方法在Kafka增加一个拦截器将数据先拦截下来先不往HDFS中存储,这个拦截器的功能是将event中bady的真实日志数据传给header,这样就修改了错误。因为HDSF是通过读取header中的时间戳对日志数据进行的分类。

三、编写拦截器

1.新建项目,相关配置如下,框起来的部分要一样

2.在pow.xml添加依赖,将原本pow.xml文件的内容全部替换为下面的

<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<groupId>com.atguigu.interceptor</groupId>
<artifactId>TimeStampInterceptor</artifactId>
<version>1.0-SNAPSHOT</version>
<properties>
<maven.compiler.source>8</maven.compiler.source>
<maven.compiler.target>8</maven.compiler.target>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
</properties>
<dependencies>
<dependency>
<groupId>org.apache.flume</groupId>
<artifactId>flume-ng-core</artifactId>
<version>1.10.1</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-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>
</project>

3.编写拦截器

创建com.atguigu.gmall.flume.interceptor包

创建TimestampInterceptor类

编写下面的代码:

package com.atguigu.gmall.flume.interceptor;

import com.alibaba.fastjson.JSONObject;
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;
import java.util.Map;

public class TimestampInterceptor implements Interceptor {

@Override
public void initialize() {

}

@Override
public Event intercept(Event event) {
//1、获取header和body的数据
Map<String, String> headers = event.getHeaders();
String log = new String(event.getBody(), StandardCharsets.UTF_8);

try {
//2、将body的数据类型转成jsonObject类型(方便获取数据)
JSONObject jsonObject = JSONObject.parseObject(log);

//3、header中timestamp时间字段替换成日志生成的时间戳(解决数据漂移问题)
String ts = jsonObject.getString("ts");
headers.put("timestamp", ts);

return event;
} catch (Exception e) {
e.printStackTrace();
return null;
}
}

@Override
public List<Event> intercept(List<Event> list) {
Iterator<Event> iterator = list.iterator();
while (iterator.hasNext()) {
Event event = iterator.next();
if (intercept(event) == null) {
iterator.remove();
}
}
return list;
}

@Override
public void close() {

}

public static class Builder implements Interceptor.Builder {
@Override
public Interceptor build() {
return new TimestampInterceptor();
}

@Override
public void configure(Context context) {
}
}
}

打成jar包,双击package

完成后会多一个target目录,如下图框起来的就是我们需要的jar包

四、编写flume日志服务脚本

编写脚本(这里的脚本与08的脚本功能不同):(hadoop102执行下面命令)

cd /home/atguigu/bin
vim f2.sh

添加下面的内容:

#!/bin/bash

case $1 in
"start")
echo " ——–启动 hadoop102 日志数据flume——-"
ssh hadoop102 "nohup /opt/module/flume/bin/flume-ng agent -n a1 -c /opt/module/flume/conf -f /opt/module/flume/job/kafka_to_hdfs_log.conf >/dev/null 2>&1 &"
;;
"stop")

echo " ——–停止 hadoop102 日志数据flume——-"
ssh hadoop102 "ps -ef | grep kafka_to_hdfs_log | grep -v grep |awk '{print \\$2}' | xargs -n1 kill"
;;
esac

添加权限:(hadoop102执行下面命令)

chmod 777 f2.sh

五、使用拦截器并将数据传输到HDFS

1.将打好的jar包放入hadoop102的/opt/module/flume/lib文件夹下(直接拖就行)

2.删除原有的日志文件

如将下图的app.log删除

3.启动需要的服务

先保证有下面的进程,如果没有输入下面命令启动服务(zookeeper和kafka)(hadoop102执行下面命令)

zk.sh start
kf.sh start

如果进程有可以输入下面命令停止(hadoop102执行下面命令)

mxw.sh stop

删除Kafka中的日志数据

如下图的topic_log和topic_db

然后再启动Hadoop集群(hadoop102执行下面命令)

hdp.sh start

查看进程,应该增加下面的进程(hadoop102执行下面命令)

xcall jps

3.启动脚本f1.sh和f2.sh将数据分别放入Kafka和HDFS(hadoop102执行下面命令)

f1.sh start
f2.sh start

4.生成日志数据(hadoop102执行下面命令)

lg.sh

六、在HDFS上查看数据

输入网址http://hadoop102:9870

就可以查看到数据了

这里说一下到目前为止的数据传递流程如图(其中的标记就是对应的启动脚本):

 

 

 

 

 

 

赞(0)
未经允许不得转载:171主机测评 » 12-学习笔记尚硅谷数仓搭建-将Kafka中的日志数据同步到Hadoop集群的HDFS分布式文件系统的flume配置
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址