实训综合案例_第1页
实训综合案例_第2页
实训综合案例_第3页
实训综合案例_第4页
实训综合案例_第5页
已阅读5页,还剩41页未读 继续免费阅读

下载本文档

版权说明:本文档由用户提供并上传,收益归属内容提供方,若内容存在侵权,请进行举报或认领

文档简介

模块七大数据日志分析综合项目案例7.1项目目的当今时代,数据与我们息息相关,我们每天都会接触到各种软件、浏览各种网站,在生活上,也不仅仅是接触信息,其实也是在生产很多信息,留下很多的数据,只是人们可能没有察觉而已。本模块将会总结前面模块所学习的大部分组件,设计成了综合项目案例,让大家对大数据的认识提到一个新的高度,并且熟悉生产上的开发流程。7.2项目意义本项目案例的日志指的是用户行为日志,用户行为日志可以类比于网站或者app的眼睛,开发人员可以从中了解到用户的主要来源、喜欢的内容、用户的访问设备等等。也可以将此比喻为网站或者app的神经,通过用户行为日志的分析,可以清楚网站或者app的优缺点,了解用户使用过程中遇到的各种问题以及反馈,进而有利于优化自己的网站或者app,提升用户的体验。此外,通过日志分析,还可以通过用户的行为日志,挖掘出有价值的信息,将信息进行归类,划分主要的倾向人群,有利于实现业务需求。7.3项目背景此处用户每次访问网站或者app时,都会留下很多的行为数据,这些行为包括访问、浏览、搜索、点击等等。每一个行为动作所产生的数据都可以被后台采集到。比如说点击的URL、从哪个URL跳转过来的(referer)、页面上的停留时间等等。当有了数据之后,就可以进行大数据的分析统计等等工作了。7.4项目架构先来了解一下本次项目的架构,再来总结数据处理的流程,如图7-1所示。图7-1项目架构流程图当访问网站或者使用app的时候,都会产生许多日志信息,存放到日志服务器里面。框架图中WebServer指的是网站或者app的后台,而本实训将日志信息直接存放到服务器的/home/access.log路径下,然后通过Flume对采集到的信息进行路由,此处路由分两条主线,一条是直接将数据采集到HDFS,让MapReduce对HDFS上的数据进行清洗或者离线分析,分析完后再将结果存放到传统数据库中,此处是使用MySQL。Flume路由的另一条主线是与Kafka整合,将消费的数据存储到由HBase中,当然此处的Kafka也可以与SparkStreaming、Storm、Flink等组件整合,实现实时流处理主线。Kafka与HBase整合完后,HBase可以与传统的业务系统整合,也可以与其他组件整合,如图7-1中将HBase与Hive进行整合,目的是实现通过类SQL对HBase中的数据进行高效的分析。最后,这两条主线可以与ECharts整合,对数据进行可视化。1.数据处理流程综上所述,可以将数据处理流程归结为五大步骤:数据采集->数据清洗->数据分析->数据入库->数据可视化1)数据采集可以使用Flume对数据进行采集,将web日志写入到HDFS、Kafka或者HBase等等中。2)数据清洗可以使用MapReduce、Spark、Hive、Flink或者其他的一些分布式计算框架,对数据进行清洗,先过滤掉没有意义的数据,如脏数据或者与业务不相关的数据等等,清洗完之后的数据可以存放在HDFS或者Hive、SparkSQL等等中。3)数据处理按照需求对相应业务进行统计和分析,可以使用数据清洗时的计算框架。4)数据处理结果入库处理的结果可以存放到RDBMS、NoSQL等数据库中。5)数据可视化当数据入库之后,可以开发各种各样的图形化界面对分析结果进行展示,比如说饼图、柱状图、地图、折线图等等,可以借助的工具有ECharts、DataV、HUE、Zeppelin、Kibana等。7.5项目需求当获取到了数据,可以从中挖掘出一些价值,想要挖掘什么价值,取决于业务能力水平,而能否实现,则取决于技术本身的能力以及所拥有的数据维度有多广有多完善,而实现的难度则与数据的质量息息相关。本次项目的数据采用模拟的方式生成,自定义的数据有ip、时间、访问的URL、跳转过来的网址、状态码。主要有5个字段,当然,此数据可以自行修改自行生成,也可以拿真实的数据来操作。基于数据,可以实现的业务场景有非常多,自己可以尝试去多挖掘。比如说,统计哪三个省份的用户访问网站最频繁?统计访问网站最频繁的时间段是哪个?统计过去10个小时内,用户的访问量有多少?还有很多,都可以实现。为了更好地与前面模块的内容衔接,也为了降低学习的难度,本次项目的业务需求是统计每天的用户访问量。7.6业务实现1.准备工作需要准备好开发工具和所需要的软件的安装包,前面的实训已经准备好了。所以此处不再做过多说明。在实操的时候,应确保各软件的版本与本书一致,不一致也应该确保大版本保持一致;如不相同,遇到问题,请先自行搜索与自己版本相关的解决方案。2.效果提前预览 项目的最终效果如图7-2所示。图7-2项目展示效果图说明:具体的次数每个人会不相同。3.实现步骤接下来将一步一步来实现,主要分为以下七大步骤:步骤一、模拟日志生产步骤二、编写Flume配置文件步骤三、Flume整合Kafka步骤四、Flume与HDFS、Kafka整合步骤五、Kafka与HBase整合步骤六、MapReduce分析HDFS上的数据并写入到MySQL步骤七、ECharts与MySQL整合实现数据可视化1)模拟日志生产①新建一个名称为logstat的项目,关键设置选项如图7-3所示。图7-3新建项目项目新建好后,界面如图7-4所示。图7-4界面总览接着,在java目录里面新建包com.bigdata.hadoop.generate操作如图7-5、图7-6所示。图7-5新建Package图7-6给新建包命名新建GenerateLog类,里面编写模拟日志生成的主程序。操作过程如图7-7、图7-8所示。图7-7新建Class图7-8给新建类命名②编写代码packagecom.bigdata.hadoop.generate;importjava.io.File;importjava.io.FileOutputStream;importjava.io.IOException;importjava.text.DateFormat;importjava.text.SimpleDateFormat;importjava.util.Calendar;importjava.util.Date;importjava.util.Random;importjava.util.concurrent.TimeUnit;publicclassGenerateLog{//一、数据定义//1、url地址publicstaticString[]urlPaths={"article/102.html","article/103.html","article/104.html","article/105.html","article/106.html","article/107.html","article/108.html","article/109.html","video/322","tag/list"};//2、ip数字publicstaticString[]ipSplices={"102","71","145","33","67","54","164","121"};//3、http网址publicstaticString[]httpReferers={"/s?wd=%s","/web?query=%s","/search?q=%s","/search?p=%s"};//4、搜索关键字publicstaticString[]searchKeyword={"复制粘贴玩大数据",

"网站用户行为分析",

"Elasticsearch的安装",

"Kafka的安装及发布订阅消息系统",

"window7系统上Centos7的安装",

"学习大数据常用Linux命令",

"Docker搭建Spark集群"};//5、状态码publicstaticString[]statusCodes={"200","404","500"};//二、随机生成数据//1、随机生成ippublicstaticStringsampleIp(){intipNum;Stringip="";for(inti=0;i<4;i++){ipNum=newRandom().nextInt(ipSplices.length);ip+="."+ipSplices[ipNum];}returnip.substring(1);}//2、随机生成时间publicstaticStringformatTime(){DateFormatdateFormat=newSimpleDateFormat("yyyy-MM-ddHH:mm:ss");Calendarcalendar=Calendar.getInstance();//获取当前时间DatecurrentDate=calendar.getTime();//设置一个起始时间(七天前)calendar.add(Calendar.DATE,-7);DatestartDate=calendar.getTime();//获取七天内的一个随机时间longdateTime=startDate.getTime()+(long)(newRandom().nextDouble()*(currentDate.getTime()-startDate.getTime()));returndateFormat.format(dateTime);}//3、随机生成urlpublicstaticStringsampleUrl(){inturlNum=newRandom().nextInt(urlPaths.length);returnurlPaths[urlNum];}//4、随机生成检索publicstaticStringsampleReferer(){Randomrandom=newRandom();intrefNum=random.nextInt(httpReferers.length);intqueryNum=random.nextInt(searchKeyword.length);if(random.nextDouble()<0.2){return"-";}Stringquery_str=searchKeyword[queryNum];StringreferQuery=String.format(httpReferers[refNum],query_str);returnreferQuery;}//5、随机生成状态码publicstaticStringsampleStatusCode(){intcodeNum=newRandom().nextInt(statusCodes.length);returnstatusCodes[codeNum];}//6、生成日志方法//输出日志格式:02,2022-11-1409:06:00,"GET/tag/listHTTP/1.1",https:///search?p=复制粘贴玩大数据,404publicstaticStringgenerateLog(){Stringip=sampleIp();StringnewTime=formatTime();Stringurl=sampleUrl();Stringreferer=sampleReferer();Stringcode=sampleStatusCode();Stringlog=ip+","+newTime+","+"\"GET/"+url+"HTTP/1.1\""+","+referer+","+code;System.out.println(log);returnlog;}//三、主类publicstaticvoidmain(String[]args)throwsIOException,InterruptedException{//dest:生成日志的路径//Stringdest="/home/access.log";Stringdest="access.log";Filefile=newFile(dest);//num:每次生成条数//sleepTime:多久生成一次intnum,sleepTime;if(args.length==2){num=Integer.valueOf(args[0]);sleepTime=Integer.valueOf(args[1]);}else{num=50;sleepTime=10;}while(true){for(inti=0;i<num;i++){Stringcontent=generateLog()+"\n";FileOutputStreamfos=newFileOutputStream(file,true);fos.write(content.getBytes());fos.close();}TimeUnit.SECONDS.sleep(sleepTime);}}}③代码解释Stringdest="/home/access.log";若解开注释,使用此行代码,则表示模拟生产的日志所存储的路径/home/access.log,需要注意的是,此为服务器上的路径,而且是root用户才具有写权限,如果不是root用户请修改成其他可写路径。在Windows系统的编辑器开发的时候,可以改成Windows上的路径来测试一下生成的日志是否为自己想要的,如D:\\access.log,表示生成的日志在D:\\access.log。若代码中直接使用access.log,表示直接生成文件在项目目录下。当执行此类时,日志就会不断生成到所设置的路径,而日志格式如图7-9所示。45,2022-11-0920:13:38,"GET/article/104.htmlHTTP/1.1",/search?q=学习大数据常用Linux命令,20064,2022-11-1011:39:16,"GET/article/104.htmlHTTP/1.1",-,20021,2022-11-1405:51:41,"GET/article/104.htmlHTTP/1.1",/search?p=Elasticsearch的安装,4043,2022-11-1208:45:44,"GET/video/322HTTP/1.1",/search?q=window7系统上Centos7的安装,4044,2022-11-1412:22:46,"GET/article/102.htmlHTTP/1.1",/web?query=Elasticsearch的安装,50064,2022-11-0804:31:54,"GET/article/104.htmlHTTP/1.1",/s?wd=网站用户行为分析,20002,2022-11-1020:13:50,"GET/article/104.htmlHTTP/1.1",/web?query=Kafka的安装及发布订阅消息系统,200图7-9日志格式生成随机时间,为了展示效果美观,本实训模拟生成七天的数据,在实际操作过程中,可以不模拟七天,直接返回当天日期即可。此外,还可以在执行的时候添加参数,第一个参数为一个批次生成的条数,第二个参数为多少秒生成一次,如不设置,则默认是每10秒生成50条。④测试生成日志接下来可以先在Windows本地测试运行,观察运行结果是否有问题。注意目前所设置的路径为:access.log。如图7-10所示。图7-10设置路径并执行点击执行按钮,稍等一小会,可以发现项目目录下有日志文件access.log生成了,如图7-11所示;控制台也有显示,如图7-12所示。图7-11查看文件日志图7-12控制台中查看结果⑤打包接下来可以将代码打包到服务器上执行,使生成的日志在服务器的/home路径下。此时需要注释掉Windows路径,修改为Linux服务器的路径,如图7-13所示。图7-13修改日志路径此外,因为服务器上的JDK版本是8,而在Windows的版本为jdk11的话,需要设置一下打包的项目语言级别才能兼容。点击“ProjectStructure”→“Project”,在“LanguageLevel”选择服务器上相应的语言级别,JDK8对应的是8级别,如图7-14所示。图7-14选择对应的语言级别此外,还需要修改一下pom.xml文件中编译代码的JDK版本,此处修改为8,默认是11。如图7-15所示。图7-15设置编译的JDK版本接着就可以将代码进行打包了,先点击编辑器右侧栏的“Maven”,再依次找到“package”,如图7-16所示。图7-16打包项目 双击“package”按钮,则可以对项目进行打包,如打包成功,控制台将会显示构建成功的标志。如图7-17所示。图7-17打包成功的标志打包完成后,发现项目里多了target文件夹,相应的jar包也生成了。如图7-18所示。图7-18查看jar包文件此时,将此jar包上传到master节点的/root/jars文件夹(没有此目录则新建创建)。如图7-19所示。图7-19查看上传路径⑥执行并查看结果(如果需要添加参数,则在后面添加上即可),任意路径执行都可以:java-cp/root/jars/logstat-1.0-SNAPSHOT.jarcom.bigdata.hadoop.generate.GenerateLog 操作结果如图7-20所示。图7-20执行模拟生成日志代码此时发现终端上一直有日志显示,此时打开一个新的终端窗口,进入到生成日志的目录查看:cd/homewc-laccess.log 操作结果如图7-21所示。图7-21查看日志行数来查看生成的日志条数,目前为700条(实操结果会有差异,日志还在实时产生)。查看前10条数据:head-n10access.log查看结果如图7-22所示。图7-22查看日志前10条数据此时先切换终端,结束之前产生日志的程序,按CTRL+C则可停止。再查看一下日志数,发现有1300条了。通过命令来查看文件的大小:du-haccess.log目前大小为156K(大概是0.12K每条日志)。操作结果如图7-23所示。图7-23查看日志行数与大小至此,模拟日志生成步骤就已经实现了。实际生产上,应该是有一个专门的服务器来生产和存储日志的,比如说日志服务器,每当用户访问Web服务器上的网站时,都会有不同的日志产生。此处不做过多介绍,如需要了解,可以自行参考Web开发相关的资料。2)编写Flume配置文件接下来,需要使用Flume来采集日志,在前面模块已经介绍过了,此处有个不一样的地方是此处不仅仅是采集到HDFS上,同时还采集到Kafka上去。①新建配置文件kafka-hdfs.conf在master节点执行:cd/opt/software/apache-flume-1.10.1-bin/confvimkafka-hdfs.conf添加内容:#agent1agent1.channels=channel1channel2agent1.sources=source1agent1.sinks=sink1sink2#agent1execSourceagent1.sources.source1.type=mand=tail-n+0-F/home/access.log#agent1memoryChannelagent1.channels.channel1.type=memoryagent1.channels.channel1.capacity=1000agent1.channels.channel1.transactionCapacity=100#agent1fileChannelagent1.channels.channel2.type=fileagent1.channels.channel2.checkpointDir=/opt/software/apache-flume-1.10.1-bin/fchannel/spool/checkpointagent1.channels.channel2.dataDirs=/opt/software/apache-flume-1.10.1-bin/fchannel/spool/dataagent1.channels.channel2.capacity=100000agent1.channels.channel2.transactionCapacity=6000agent1.channels.channel2.checkpointInterval=60000#agent1hdfsSinkagent1.sinks.sink1.type=hdfsagent1.sinks.sink1.hdfs.path=hdfs://master:8020/user/flume/events/%Y-%m-%dagent1.sinks.sink1.hdfs.filePrefix=eventsagent1.sinks.sink1.hdfs.rollInterval=600agent1.sinks.sink1.hdfs.rollSize=268435456agent1.sinks.sink1.hdfs.rollCount=0agent1.sinks.sink1.hdfs.idleTimeout=3600agent1.sinks.sink1.hdfs.writeFormat=Textagent1.sinks.sink1.hdfs.inUseSuffix=.txtagent1.sinks.sink1.hdfs.fileType=DataStreamagent1.sinks.sink1.hdfs.useLocalTimeStamp=true#agent1kafkaSinkagent1.sinks.sink2.type=org.apache.flume.sink.kafka.KafkaSinkagent1.sinks.sink2.topic=kafkatopicagent1.sinks.sink2.brokerList=master:9092,slave1:9092,slave2:9092agent1.sinks.sink2.requiredAcks=1agent1.sinks.sink2.batchSize=20#将source和sink绑定到channelagent1.sources.source1.channels=channel1channel2agent1.sinks.sink1.channel=channel2agent1.sinks.sink2.channel=channel1 此处额外添加了一个fileChannel和kafkaSink。写好配置文件之后,可以先不启动,等后面的组件整合完成再联调也可以。3)Flume整合Kafka注意到上面Flume的配置文件里,Kafka的Topic取名为kafkatopic,所以在Kakfa里需要新建一个kafkatopic的Topic。①新建Kafka的Topic启动ZooKeeper(三台服务器都要执行):zkServer.shstart启动Kafka(三台服务器都要执行):kafka-server-start.sh-daemon$KAFKA_HOME/config/perties此时查看一下各节点的进程情况。如图7-24所示。~/shell/jps_all.sh图7-24查看三台节点的进程情况接下来新建名称为kafkatopic的Topic:kafka-topics.sh--create--replication-factor3--partitions5--topickafkatopic--bootstrap-servermaster:9092,slave1:9092,slave2:9092 创建结果如图7-25所示。图7-25新建Topic创建好后,可以查看一下Topic的详情,观察是否正常:kafka-topics.sh--describe--topickafkatopic--bootstrap-servermaster:9092,slave1:9092,slave2:9092如图7-26所示,表示节点正常。图7-26查看Topic详情②构建Kafka业务代码结构继续回到IDEA编辑器里编写代码,因为生产者来源于Flume。所以,此处不需要自己编写生产者,但是在测试的时候,自己应该养成编写测试类的习惯,编写生产者来调试,模拟生产一些数据,此处省略调试过程,只实现了消费者端。新建包名(注意位置和包名)如图7-27所示。图7-27新建包如图7-28所示,新建CustomConsumer类:图7-28新建CustomConsumer类③引入Kafka所需要的pom.xml依赖因为项目里要用到打包的相关插件,所以一起将其加进来,加粗字体为新增的内容。目前完整的pom.xml文件参考如下:<?xmlversion="1.0"encoding="UTF-8"?><projectxmlns="/POM/4.0.0"xmlns:xsi="/2001/XMLSchema-instance"xsi:schemaLocation="/POM/4.0.0/xsd/maven-4.0.0.xsd"><modelVersion>4.0.0</modelVersion><groupId>com.bigdata.hadoop</groupId><artifactId>logstat</artifactId><packaging>pom</packaging><version>1.0-SNAPSHOT</version><properties><piler.source>8</piler.source><piler.target>8</piler.target><project.build.sourceEncoding>UTF-8</project.build.sourceEncoding><kafka.version>3.3.1</kafka.version></properties><dependencies><dependency><groupId>org.apache.kafka</groupId><artifactId>kafka-clients</artifactId><version>${kafka.version}</version></dependency><dependency><groupId>org.slf4j</groupId><artifactId>slf4j-nop</artifactId><version>1.7.36</version></dependency></dependencies><build><plugins><plugin><groupId>org.apache.maven.plugins</groupId><artifactId>maven-compiler-plugin</artifactId><version>3.8.0</version><configuration><source>1.8</source><target>1.8</target><testExcludes><testExclude>/src/test/**</testExclude></testExcludes><encoding>utf-8</encoding></configuration></plugin><plugin><artifactId>maven-assembly-plugin</artifactId><configuration><descriptorRefs><descriptorRef>jar-with-dependencies</descriptorRef></descriptorRefs></configuration><executions><execution><id>make-assembly</id><!--thisisusedforinheritancemerges--><phase>package</phase><!--指定在打包节点执行jar包合并操作--><goals><goal>single</goal></goals></execution></executions></plugin></plugins></build></project> 添加好后,需要导入依赖。右击pom.xml文件,选择“Maven”,选择“Reloadproject”。如图7-29所示。图7-29导入kafka依赖kafka-clients是进行Kafka编程所需要导入的依赖,slf4j-nop是为了解决运行时报警告而引入的依赖。等加载完,则可以继续操作。④定义配置项由于编写代码过程中会用到相关的配置类,此时可以定义一个专门类来存放,在hadoop包下新建property包,并且在此包中新建MyProperties类。如图7-30所示。图7-30新建MyProperties类在MyProperties类中添加配置项:packageperty;publicclassMyProperties{//Kafka相关配置项publicstaticfinalStringZK="31:2181";publicstaticfinalStringTOPIC="kafkatopic";publicstaticfinalStringBROKER_SERVER="31:9092";publicstaticfinalStringGROUP_ID="group1";}⑤编写CustomConsumer类代码packagecom.bigdata.hadoop.kafka;importcom.bigdata.hadoop.hbase.HBaseDAO;importperty.MyProperties;importorg.apache.kafka.clients.consumer.ConsumerRecord;importorg.apache.kafka.clients.consumer.ConsumerRecords;importorg.apache.kafka.clients.consumer.KafkaConsumer;importjava.time.Duration;importjava.util.Arrays;importjava.util.Properties;publicclassCustomConsumer{publicstaticvoidmain(String[]args){Propertiesprops=newProperties();//Kafka集群props.put("bootstrap.servers",MyProperties.BROKER_SERVER);//消费者组,只要group.id相同,就属于同一个消费者组props.put("group.id",MyProperties.GROUP_ID);//关闭自动提交offsetprops.put("mit","false");//设置key和value的反序列化方式props.put("key.deserializer","mon.serialization.StringDeserializer");props.put("value.deserializer","mon.serialization.StringDeserializer");KafkaConsumer<String,String>consumer=newKafkaConsumer<>(props);//消费者订阅主题consumer.subscribe(Arrays.asList(MyProperties.TOPIC));while(true){//消费者拉取数据ConsumerRecords<String,String>records=consumer.poll(Duration.ofMillis(100));for(ConsumerRecord<String,String>record:records){//System.out.println("offset=%d,key=%s,value=%s%n",record.offset(),record.key(),record.value());Stringmessage=record.value();System.out.println("Receive:"+message);}//同步提交,当前线程会阻塞直到offset提交成功mitSync();}}}⑥打包项目到集群并执行: 重新打包Maven项目,稍等片刻,等打包完会发现多了一个以-jar-with-dependencies结尾的jar包,将此jar包上传到master服务器的~/jars路径,并执行:java-cp/root/jars/logstat-1.0-SNAPSHOT-jar-with-dependencies.jarcom.bigdata.hadoop.kafka.CustomConsumer可以发现没有数据输出,因为Flume还没有启动。 操作结果如图7-31所示。图7-31执行项目代码此时再切换终端,查看发现多了一个CustomConsumer进程,如图7-32所示。图7-32切换终端并查看进程4)Flume与HDFS、Kafka整合①先启动HDFS和YARN:start-all.sh启动完成后查看各节点进程,如图7-33所示则表示各服务都正常。~/shell/jps_all.sh图7-33查看三台节点的进程情况②启动Flume(注意目前的执行路径为$FLUME_HOME/conf):flume-ngagent--conf$FLUME_HOME/conf--conf-file$FLUME_HOME/conf/kafka-hdfs.conf--nameagent1Dflume.root.logger=DEBUG,console 操作结果如图7-34所示。图7-34启动Flume打开一个新的终端窗口,查看HDFS上的数据,发现Flume的数据已经生产到了HDFS上了(配置文件里配置的HDFS路径为/user/flume/events)hdfsdfs-ls/user/flume/events/ 操作结果如图7-35所示。图7-35查看HDFS上是否有数据查看处于执行状态的CustomConsumer程序终端,也有数据显示。如图7-36所示。图7-36启动CustomConsumer程序的终端情况此时,将生成日志的程序启动:java-cp/root/jars/logstat-1.0-SNAPSHOT.jarcom.bigdata.hadoop.generate.GenerateLog启动后,再观察启动CustomConsumer程序的终端窗口,其实也是会有数据不断打印出来的。至此,Flume、Kafka、HDFS就已经整合好了。5)Kafka与HBase整合①引入HBase所需要的pom.xml依赖(注意位置)<hbase.version>2.5.0</hbase.version><dependency> <groupId>org.apache.hbase</groupId> <artifactId>hbase-client</artifactId> <version>${hbase.version}</version></dependency> 操作结果如图7-37所示。图7-37引入HBase相关依赖②在MyProperties配置类中添加编写HBase代码相关配置//HBase相关配置项publicstaticfinalStringZK_NODE="/hbase";publicstaticfinalStringTABLENAME="loginfo";publicstaticfinalIntegerPARTITION_NUM=100;③编写HBaseDAO类代码在hadoop包下新建hbase包,新建HBaseDAO类,结构如图7-38所示。图7-38新建HBase包和HBaseDAO类HBaseDAO类完整代码如下:packagecom.bigdata.hadoop.hbase;importperty.MyProperties;importorg.apache.hadoop.conf.Configuration;importorg.apache.hadoop.hbase.*;importorg.apache.hadoop.hbase.client.*;importjava.io.IOException;import.MalformedURLException;import.URL;importjava.text.DecimalFormat;publicclassHBaseDAO{privateTabletable=null;privateTableNametableName=null;//HBase的分区个数privateintpartitonsNum=0;//一、初始化publicHBaseDAO(){//1、配置项Configurationconfiguration=HBaseConfiguration.create();Connectionconnection=null;configuration.set("hbase.zookeeper.quorum",MyProperties.ZK);configuration.set("zookeeper.znode.parent",MyProperties.ZK_NODE);partitonsNum=MyProperties.PARTITION_NUM;try{//2、获取连接connection=ConnectionFactory.createConnection(configuration);//3、获取HBaseAdmin对象Adminadmin=connection.getAdmin();Stringtbl=MyProperties.TABLENAME;TableNametableName=TableName.valueOf(tbl);//4、表不存在时创建表if(!admin.tableExists(tableName)){//创建表描述对象TableDescriptorBuildertableDescriptor=TableDescriptorBuilder.newBuilder(tableName);//列簇1ColumnFamilyDescriptorfamilyColumn1=ColumnFamilyDescriptorBuilder.newBuilder("c1".getBytes()).build();//列簇2ColumnFamilyDescriptorfamilyColumn2=ColumnFamilyDescriptorBuilder.newBuilder("c2".getBytes()).build();tableDescriptor.setColumnFamily(familyColumn1);tableDescriptor.setColumnFamily(familyColumn2);//用HBaseAdmin对象创建表admin.createTable(tableDescriptor.build());}//5、获取表table=connection.getTable(tableName);//6、关闭HBaseAdmin对象admin.close();}catch(IOExceptione){e.printStackTrace();}}//二、向表put数据//日志初始格式:02,2022-11-1409:06:00,"GET/tag/listHTTP/1.1",/search?p=复制粘贴玩大数据,404//想要输出的结果:02,20221114090600,/tag/list,,404//Stringlog=ip+"\t"+newTime+"\t"+"\"GET/"+url+"HTTP/1.1\""+"\t"+referer+"\t"+code;//ip:02//date:2022-11-1409:06:00//actionSource:"GET/tag/listHTTP/1.1"//refererSource:/search?p=复制粘贴玩大数据//code:404publicvoidput(Stringlog){//1、切割日志信息String[]arr=log.split(",");Stringip=arr[0];Stringdate=arr[1];StringactionSource=arr[2];StringrefererSource=arr[3];Stringcode=arr[4];System.out.println("原始数据:"+ip+","+date+","+actionSource+","+refererSource+","+code);//2、转换格式为以获取的想要的结果StringdateFormat=date.replace("-","").replace(":","").replace("","");//删除“-”、空格和“:”String[]actionArr=actionSource.split("");//删除空格Stringaction=actionArr[1];Stringreferer=null;if(!refererSource.equals("-")){String[]refererArr=refererSource.split("\\?");try{URLurl=newURL(refererArr[0]);referer=url.getHost();}catch(MalformedURLExceptione){e.printStackTrace();}}else{referer="-";}//3、计算出该日志所在的Region区域号StringhashCode=getHashCode(ip,dateFormat);//4、拼接HBase的RowKeyStringrowKey=hashCode+","+ip+","+dateFormat+","+code;//测试输出结果的结果:02,20221114090600,/tag/list,,404System.out.println("转化后数据:"+ip+","+dateFormat+","+action+","+referer+","+code);//5、创建Put对象Putput=newPut(rowKey.getBytes());//在列簇中添加相应的列put.addColumn("c1".getBytes(),"ip".getBytes(),ip.getBytes());put.addColumn("c1".getBytes(),"date".getBytes(),dateFormat.getBytes());put.addColumn("c1".getBytes(),"action".getBytes(),action.getBytes());put.addColumn("c1".getBytes(),"referer".getBytes(),referer.getBytes());put.addColumn("c1".getBytes(),"code".getBytes(),code.getBytes());//put数据到表中try{table.put(put);}catch(IOExceptione){e.printStackTrace();}}//三、计算出该日志所在的Region区域号实现方法privateStringgetHashCode(Stringip,StringdateFormat){//1、此处取ip地址最后5位StringipNum=ip.replace(".","");intlen=ipNum.length();Stringnumber=ipNum.substring(len-4);//2、取出年份和月份Stringdate=dateFormat.substring(0,6);//3、随机数取哈希值intcode=(Integer.parseInt(date)^Integer.parseInt(number))%partitonsNum;//4、格式化后返回DecimalFormatdf=newDecimalFormat();df.applyPattern("00");returndf.format(code);}}④测试代码在hbase包下简单写一个HBaseDAOTest类,在里面编写一个main方法,调用HBaseDAO执行插入一条日志的操作。完整代码如下:packagecom.bigdata.hadoop.hbase;publicclassHBaseDAOTest{publicstaticvoidmain(String[]args){//定义一条数据Stringlog="02,2022-11-1409:06:00,\"GET/tag/listHTTP/1.1\",/search?p=复制粘贴玩大数据,404";//插入数据测试HBaseDAOhBaseDAO=newHBaseDAO();hBaseDAO.put(log);}}启动HBase集群:start-hbase.sh执行main方法,执行结束后进入HBaseShell操作页面,可用查看到已经通过代码新建了表,如图7-39所示。hbaseshelllist图7-39查看表查看loginfo表数据,如图7-40所示。scan"loginfo"图7-40查看表数据可以看到,已经将日志数据插入到了HBase表中,说明测试成功。测试成功后,我们需要删除一下测试数据,先disable表,再drop表。如图7-41所示。disable'loginfo'drop'loginfo'图7-41删除表⑤CustomConsumer类调用HBaseDAO测试插入数据到HBase成功后,此时就可以将Kafka与HBase结合起来了。所以,此时可以在Kafka的消费者端调用HBaseDAO相关代码,直接在CustomConsumer的main方法里调用即可,在while循环里调用put方法将数据写入到HBase中,代码如图7-42、7-43所示。HBaseDAOhbaseDao=newHBaseDAO();hbaseDao.put(message);图7-42构建HBaseDAO对象图7-43CustomConsumer类调用HBaseDAO⑥打包到集群并执行重新打包新程序并上传到服务器上,执行CustomConsumer程序。java-cp/root/jars/logstat-1.0-SNAPSHOT-jar-with-dependencies.jarcom.bigdata.hadoop.kafka.CustomConsumer切换终端,查看HBase的表,可以看到生成了loginfo表,但此时还没有数据。如图7-44所示。图7-44查看HBase的表此时需要将Flume和模拟生成日志程序启动好。打开一个新终端,执行:flume-ngagent--conf$FLUME_HOME/conf--conf-file$FLUME_HOME/conf/kafka-hdfs.conf--nameagent1Dflume.root.logger=DEBUG,console打开一个新终端,启动模拟生成日志程序:java-cp/root/jars/logstat-1.0-SNAPSHOT.jarcom.bigdata.hadoop.generate.GenerateLog查看表数据:scan"loginfo" 查看结果如图7-45所示。图7-45查看HBase表信息并且可以发现,loginfo表中的数据,是一直在增加的。因为是实时插入的。 至此,模拟日志生成,Flume实时采集日志,将采集的日志发送给Kafka,并且将数据写到到HBase的流程就跑通了。6)MapReduce分析HDFS上的数据并写入到MySQL接下来就要对数据进行分析了,为了简化操作流程,本项目不考虑其他因素,直接将IP作为用户的访问量,实际开发上还要考虑很多因素的,而且也已经将数据模拟得非常工整。也就是说,默认一个IP就是一个用户,前面已经提及,本项目业务需求是统计每天用户的访问量,需要将统计结果写入到MySQL里。此过程对学生所掌握的基础要求比较高,知识点比较多,如果接触SpringBoot和JS、HTML的话会比较容易上手。①MySQL准备工作打开一个新终端,先登录MySQL:mysql-uroot-p123456新建数据库logstat:createdatabaselogstat;uselogstat;创建统计结果相应的表day_log_access_topn_stat:createtableday_log_access_topn_stat(dayvarchar(8)notnull,timesbigint(10)notnull,primarykey(day));②引入HDFS所需要的pom.xml依赖(注意位置)<hadoop.version>3.3.4</hadoop.version><dependency><groupId>org.apache.hadoop</groupId><artifactId>hadoop-client</artifactId><version>${hadoop.version}</version></dependency> 加入依赖位置如图7-46所示。图7-46引入编写Hadoop程序所需的依赖编写好pom.xml文件后,记得Reload一下项目,将依赖加载到项目中。③编写代码准备好后,新建mapreduce包和LogStat2MySQL类,如图7-47所示。图7-47新建包和类在MyProperties类中加入MySQL配置项。//MySQL相关配置项publicstaticfinalStringDRIVER="com.mysql.cj.jdbc.Driver";publicstaticfinalStringURL="jdbc:mysql://master:3306/logstat?useUnicode=true&characterEncoding=UTF8";publicstaticfinalStringUSERNAME="root";publicstaticfinalStringPASSWORD="123456";编写LogStat2MySQL类packagecom.bigdata.hadoop.mapreduce;importjava.io.DataInput;importjava.io.DataOutput;importjava.io.IOException;importjava.sql.PreparedStatement;importjava.sql.ResultSet;importjava.sql.SQLException;importperty.MyProperties;importorg.apache.hadoop.io.Writable;importorg.apache.hadoop.mapred.lib.db.DBWritable;importorg.apache.hadoop.conf.Configuration;importorg.apache.hadoop.fs.FileSystem;importorg.apache.hadoop.fs.Path;importorg.apache.hadoop.io.LongWritable;importorg.apache.hadoop.io.Text;importorg.apache.hadoop.mapreduce.Job;importorg.apache.hadoop.mapreduce.Mapper;importorg.apache.hadoop.mapreduce.Reducer;importorg.apache.hadoop.mapreduce.lib.db.DBConfiguration;importorg.apache.hadoop.mapreduce.lib.db.DBOutputFormat;importorg.apache.hadoop.mapreduce.lib.input.FileInputFormat;importorg.apache.hadoop.mapreduce.lib.output.FileOutputFormat;publicclassLogStat2MySQL{//一、与MySQL整合publicstaticclassTblsWritableimplementsWritable,DBWritable{Stringday;inttimes;publicTblsWritable(){}publicTblsWritable(Stringday,inttimes){this.day=day;this.times=times;}publicvoidwrite(PreparedStatementstatement)throwsSQLException{statement.setString(1,this.day);statement.setInt(2,this.times);}publicvoidreadFields(ResultSetresultSet)throwsSQLException{this.day=resultSet.getString(1);this.times=resultSet.getInt(2);}publicvoidwrit

温馨提示

  • 1. 本站所有资源如无特殊说明,都需要本地电脑安装OFFICE2007和PDF阅读器。图纸软件为CAD,CAXA,PROE,UG,SolidWorks等.压缩文件请下载最新的WinRAR软件解压。
  • 2. 本站的文档不包含任何第三方提供的附件图纸等,如果需要附件,请联系上传者。文件的所有权益归上传用户所有。
  • 3. 本站RAR压缩包中若带图纸,网页内容里面会有图纸预览,若没有图纸预览就没有图纸。
  • 4. 未经权益所有人同意不得将文件中的内容挪作商业或盈利用途。
  • 5. 人人文库网仅提供信息存储空间,仅对用户上传内容的表现方式做保护处理,对用户上传分享的文档内容本身不做任何修改或编辑,并不能对任何下载内容负责。
  • 6. 下载文件中如有侵权或不适当内容,请与我们联系,我们立即纠正。
  • 7. 本站不保证下载资源的准确性、安全性和完整性, 同时也不承担用户因使用这些下载资源对自己和他人造成任何形式的伤害或损失。

评论

0/150

提交评论