hadoop6天3笔记1 1大数据部分课程介绍_第1页
hadoop6天3笔记1 1大数据部分课程介绍_第2页
hadoop6天3笔记1 1大数据部分课程介绍_第3页
hadoop6天3笔记1 1大数据部分课程介绍_第4页
hadoop6天3笔记1 1大数据部分课程介绍_第5页
已阅读5页,还剩201页未读 继续免费阅读

下载本文档

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

文档简介

HADOOP 学习建 HADOOP简 前 hadoop应用场 hadoop集群部署安 hdfs的s操 HDFS的java操 hdfs的工作机 namenode工作机 datanode的工作机 一些补 hdfs读数据流 hdfs写数据流 hadoop的RPC框 hdfs读数据源码分 hdfs写数据源码分 debugHadoop服务端代码 MAPREDUCE编程规 MAPREDUCE中的 Mapreduce的排序初 Partitioner编 Mapreduce的排 partital排序示例,多reducetask自动实现各输出文件有 total排序机 secondary排序机 shuffle详 mr程序map任务数的规划机 Mapreduce的join算 mapreduce的Distributed Mapreduce输入格式组 由maptask数量的决定机制引入 InputFormat的继承体 自定义 Mapreduce输出格式组 自定义 Configuration配置对象与 mapreduce数据压 mapreduce的计数 mapreduce的日志分 多job串 MapReduce程序向yarn提交执行的流程分 资源请 任务调度--capacityscheduler/fair Scheduler概 CapacityScheduler配 FairScheduler配 hadoop-sp问题及HA解决思 zookeeper简 zookeeper集群搭 zookeeper演示测 zookeeper-api应 demo增删改 CDH介绍,演 hbase数据库介 hbase集群结 hbase集群搭 基本s命 hbase工作原 Region管 Master工作机 hbase应用案例看行键设 Hbase和mapreduce结 从Hbase中数据写入 从Hbase中数据写入 hbase高级编 协处理 二级索引 Hive基 hive引 hive技术架 Hive的安装部 Hive使用方 hql基本语 基本hql语 hql查询进 hive数据类 Hive高级应 Hive常用函 hive高级操 hive优 总 工 flume介 Storm基 storm介 storm基本概 storm集群搭 storm示例编 kafka介绍与应用开 Storm高级特 Storm与kafka整 Storm的tuple storm事务 引 StormTopology的ack机 seedertuple的状态 机器学 spark框架介 spark集群概 spark优势简 spark集群搭 spark编程基 spark-s编 ide编 spark-RDD原理深度解 spark内核源码阅 akka介绍,demo示例 RDD SparkContext stage划 task提 hdfs flume socket window操 statefulwindow操 sparksql概 sparksql开发示 MLlib介 HADOOPHadoop(hdfs、mapreduce、yarn) HivesqlMRSqoopFlume框StormSparkallinone,新秀,发展势头迅猛 其他公 薪 --高级开发人员平台开发--架构级别| 架构 管HADOOPjava程序hadoophadoop(比如structsspring从另一个角度,hadoop又可以理解为一个提供服务的软件(比如数据库服务具体来说,狭义上的hadoop两个大的功能:海量数据的;海量数据的分析(编程Hadoop有3大组件hadoop分布式文件系统海量数据的(集群服务),Hadoop最早来自于的三 (为什么会需要这么一种技术后来经过dougcutting的山寨,出现了java版本的 mapreduce和以上三个组件整合起来成为apache的一个顶级项目一名,从此,hadoop声名鹊起,风靡全球0.20.2--1.2.1→|→2.2.0(HDFS的namenode高可用)--→2.4.1(YARN的高可用)-- 经过演化,hadoopyarn(mapreduceyarn1linux服务器(centos6.432位<64位,可以支持更大的解压centos虚拟机镜像压缩包到某个 并用vmware打开(icopyit,起作用的网卡会变成eth1)准备操作系统环境(主机名,ip地址配成static,和ip的本地映射hosts)0ip地址setup1、更改主机名(必须用root) 2、内网映射配置 scp/etc/hostshdp-node02:/etc/scp/etc/hostshdp- 配置(关闭) serviceiptables chkconfigiptables为hadoop软件准备一个专门的普通用户(不要用root直接安装软件)可选:为hadoop用户设置sudo权限 visftpputc解压安装包 /root/ - 修改环境变量:viexportexport应该把jdk的安 和exportexportsource3、安装 #4、修改配置文件(参考现成的配置文件xxx-site.xml) 指定hadoop yarn.resourcemanager.hostname:server01 ,在hdfs-site.xml ,在hdfs-site.xmlsecondarynamenode有些公司会采用一些商业版(CDH--cloudera公司的产品;HORTONWORKS;5在相应服务器上启动hdfsnamenodesbin/hadoop-daemon.shstartdatanodesbin/hadoop-daemon.shstartbin/hdfsdfsadmin hdfs服务:sbin/start-dfs.shyarn服务:sbin/start-yarn.shhdfs+yarnsbin/start-加入文件~/.ssh/authorized_keys (该文件的权限:600)ssh工具箱的工具:1/在登陆方生成密钥对,执行命令:ssh-keygen hdfs文件系统会给客户端提供一个统一的抽象 树,客户端hdfs文件时就是通 Hdfs中的文件都是分块(block)的,块的大小可以通过配置参数(来规定,默认大小在hadoop2.x128M Hdfs中有一个重要的角色:namenode,负责整个hdfs文件系统的 一个路径(文件)所对应的block块信息(block的id,及所在的datanode服务器)[-appendToFile[-appendToFile<localsrc>...<dst>][-cat[-ignoreCrc]<src>...][-checksum<src>[-chgrp[-R]GROUP od[-R]<MODE[,MODE]...|OCTALMODE>[-chown[-R][OWNER][:[GROUP]][-copyFromLocal[-f][-p]<localsrc>...[-copyToLocal[-p][-ignoreCrc][-crc]<src>...<localdst>][-count[-q]<path>...][-cp[-f][-p]<src>...[-createSnapshot<snapshotDir>[<snapshotName>]][-deleteSnapshot<snapshotDir><snapshotName>][-df[-h][<path>[-du[-s][-h]<path>...][-get[-p][-ignoreCrc][-crc]<src>...<localdst>][-getfacl[-R]<path>][-getmerge[-nl]<src><localdst>][-help[cmd...]][-ls[-d][-h][-R][<path>[-mkdir[-p]<path>[-moveFromLocal<localsrc>...<dst>][-moveToLocal<src><localdst>][-mv<src>...[-put[-f][-p]<localsrc>...[-renameSnapshot<snapshotDir><oldName><newName>][-rm[-f][-r|-R][-skipTrash]<src>...][-rmdir[--ignore-fail-on-non-empty]<dir>[-setfacl[-R][{-b|-k}{-m|-x<acl_spec>}<path>]|[--set<acl_spec><path>]][-setrep[-R][-w]<rep><path>...][-stat[format]<path>[-tail[-f][-test-[defsz][-text[-ignoreCrc]<src>...][-touchz<path>...][-usage[cmd- #显 信-->hadoopfs-lshdfs://hadoop-这些参数中,所有的hdfs-->hadoopfsls - #hdfs-->hadoopfs-mkdir-p- - #从hdfs-- -- - - -->hadoopfs- -->hadoopfs-od666/- -- - Eg:hadoopfs-copyToLocal- -->hadoopfs-count- #从hdfshdfshadoopfs- 信息快-->hadoopfs-createSnapshot- >hadoopfs-df-h-->hadoopfs-du-s-h- - -->比如hdfs的 /aaa/下有多个文件:log.1,log.2,log.3,...hadoopfs-getmerge/aaa/log.*./log.sum- - - - -->hadoopfs-rm-r- #- -->hadoopfs-setrep3- - - hdfsDataNodeHDFS1namenode通信请求上传文件,namenode检查目标文件是否已存在,父3、请求第一个block该传输到哪些datanode服务器4、namenode3datanode5、请求3台dn中的一台A上传数据(本质上是一个RPC调用,建立BBC,将真个pipeline、 、7、当一个block传输完成之后,再次请求namenode上传第二个block的服务器HDFS1namenode通信查询元数据,找到文件块所在的datanode2、挑选一台datanode(就近原则,然后随机)socket4、客户端以packetnamenodenamenodehdfs元数据是怎么的Chdfsedits这log日志中,当客户端操作成功后,相应的元数据会更新到内存中每隔一段时间,会由secondarynamenode将namenode上积累的所有edits和一个fsimage到本地,并加载到内存进行merge(这个过程称为checkpoint)D、checkpoint操作的触发条件配置参数: #检查触发条件是否满足的频率,60#checkpoint操作时,secondarynamenode的本地工作 # #检查触发条件是否满足的频率,60#checkpoint操作时,secondarynamenode的本地工作 #两次checkpoint之间的时间间隔3600秒 #两次checkpoint之间最大的操作记录Fhdfseditsbin/hdfsoev-iedits-odatanodeDatanodenamenodeblock信息(通过心跳信息上报)block具体的物理存放情况 HDFSjava 搭建开发环境(eclipse,hdfs的jar包hadoop的安 如果非要在windowA、在windows的某 下解压一个hadoop的安装B、将安装包下的lib和 用对应windows版本平台编译的本地库替DwindowspathhadoopConfigurationconf=newConfigurationconf=newFileSystemfs=而我们的操作目标是HDFS,所以获取到的fsDistributedFileSystem从conffs.defaultFSclasspath下也没有给定相应的配置,conf中的默hadoopjar包中的core-default.xml,默认值为:file:///fs(3)hdfsHDFS数据高可靠数据延迟较大,不支持数据的修改操作适合一次写入多次的应用场景1hdfs2、对hdfs运作的理解 文件上传、的流 HDFS的其他方式HDFShdfssrestapijavaapifuse这种工HDFS还可以挂载为一个NFSFileUtilFileUtil.copy(newFile(c:/test.tar.gz),FileSystem.get(URI.create(hdfs://hadoop-server01:9000),conf,hadoop),newPath(/test.tar.gz),true,HDFStrash Namenodenameonde发现文件block丢失的数量达到一个配置阈值时,就会进入安全模式,datanode向它汇报block信息。在安全模式下,namenode可以提供元数据查询的功能,但是不能修改;namenode的安全模式:hdfsdfsadmin- <enter|leave|get|2.7.5可以随机定位位置hadoopRPCHadoop中各节点之间存在大量的过程调用,hadoop为此封装了一个RPC基础框架public publicstaticfinallongversionID=publicStringgetMetaData(String}@authorpublicclassNamNodeNameSystemImplNameNodeProtocalpublicStringgetMetaData(Stringpath)//manylogiccodetofindthemetadatainmetadatapoolreturn"{/aa/bb/bian4.mp4;300M;[BLK_1,BLK_2,BLK_3];3;}}RCP@authorRPCRCP@authorpublicclassPublishServiceToolpublicstaticvoidmain(String[]args)throwsHadoopIllegalArgumentException,{Builderbuilder=newRPC.Builder(news).setInstance(newServerserver=}}publicpublic{publicstaticvoidmain(String[]args)throwsException (NameNodeProtocal.class,1L,newnew 业务类的实现方法(本质上,具体实现在远端,走的是socket通信请求)StringmetaData=namenodeImpl.getMetaData("/aa/bb/bian4.mp4");}}(4)RPC框架APIRPChdfshdfsdebugHadoop服务端代码(扩展)# # 在本地的eclipse中打开NameNode或者DataNode类,点击右键,添加debug配添加一个debug调试配填写服务端的debug地址和端接着在namenode回到eclipse之前配置的debug配置上,点击debug开始调MAPREDUCE,Mapreduce是一个分布式运算的编程框架功能是将用户编写的逻辑代码分布式,学:掌握MR程序编程规范;MRMRmapreduce框架后,开发人员可以将绝大部分工作集中在业务逻辑的开发上,而2.4.1.jarhdfs,yarn然后在集群中的任意一台服务器上执行,(比如运行hadoopjarhadoop-mapreduce-example-2.4.1.jar /wordcount/dataMAPREDUCEMapperKV对的形式,KVMapperKV对的形式,KVMappermapmap方法是(maptask进程)每进来一个KVReducer的业务逻辑写在reduceMapperReducer整个程序需要一个Drvierjobwordcountmapper valuein://keyout: //map//key //value:protectedvoidmap(LongWritablekey,Textvalue,Contextcontext)throwsIOException,InterruptedException{Stringline=Stringwordsline.splitfor(Stringword:words){}}}//kv组,//kv组,reduceprotectedvoidreduce(Textkey,I ble<IntWritable>values,Contextcontext)throwsIOException,InterruptedException{intcount=0;//kvv,累加到countfor(IntWritablevalue:values){count+=value.get();}}}的结果放哪里。。。。。。)job对象publicpublicstaticvoidmain(String[]args)throwsException{Configurationconf=newConfiguration();Jobwcjob= booleanres=}MAPREDUCEdebug怎样实现本地运行?:写一个程序,不要带集群的配置文件(mrdebug,只要在eclipse$hadoopjarwordcountjar Blinux的eclipsemainmaee.raewr.ae=ynC、如果要在windowseclipsejobYarnRunnermaptaskmaptask数量的决定机制——2、分配的就是将原始数据进行“切片”,每一片就交给一个maptask来处比如:file1.txt mapreduce.input.fileinputformat.split.maxsizeblocksize小mapreduce.input.fileinputformat.split.minsizeblockSize大,则可以让切blocksize还大Reducetask数量的决定是由代码中api//11个reducetaskcombinerMRMapperReducercombinercombiner和reducer的区别在于运行的位置:Combinermaptask所在的节点运行ReducerMapper的输出结果;combinermaptask的输出进行局部汇总,以减小网络传输量1combinerReducerreduce2、在job中设置 而且,combiner的输出kvreducer的输入kv(1)Java的序列化是一个重量级序列化框架(Serializable),一个对象被序列化后,会附publicclassTestSeripublicclassTestSeripublicstaticvoidmain(String[]args)throwsExceptionByteArrayOutputStreamba=newByteArrayOutputStream();ByteArrayOutputStreamba2=newByteArrayOutputStream();//DataOutputStreamjdk标准序列化DataOutputStreamdout=newDataOutputStream(ba);DataOutputStreamdout2=newDataOutputStream(ba2);ObjectOutputStreamobout=newObjectOutputStream(dout2);ItemBeanSeritemBeanSer=newItemBeanSer(1000L,89.9f);ItemBeanitemBean=newItemBean(1000L,89.9f);Textatext=new//atext.write(dout);byte[]byteArray=for(byteb:byteArray){} Stringastr=//dout2.writeUTF(astr);byte[]byteArray2=ba2.toByteArray();for(byteb:byteArray2){}}}如果需要将自定义的bean放在keycomparablepublicclassFlowBean **反序列化的方法,反序列化时,从流中publicvoidreadFields(DataInputin)throwsIOExceptionupflow=in.readLong();dflow=in.readLong();sumflow=in.readLong();}*publicvoidwrite(DataOutputout)throwsIOException 的 }publicintcompareTo(FlowBeano)}MapreduceMR程序在处理数据的过程中会对数据排序(mapkvreduce之前,会排序),mapkey key的compareToMapreduceMapreducemapkvkeymapreduce之间的数据(分组)mapreduce之间的数据(分组)**publicclassProvincePartitionerextendsPartitioner<Text,FlowBean>staticHashMap<String,Integer>provinceMap=newHashMap<String,Integer>();static{}publicintgetPartition(Textkey,FlowBeanvalue,intnumPartitions){Integercode=provinceMap.get(key.toString().substring(0,3));returncode==null?5:}}map方法中,可以从参数context中获取到当前所处理的行所在的切片信息 (FileSplit)Context.getSplit();StringfileName=Split.getPath().getName()mapreduceshuffleshuffle:洗牌、发牌——(机制:数据分区,排序,缓存6reducetaskmaptask的结果文件,reducetask会将这些(group,调用用户自定义的reduce()方法Shufflemapreduce程序的执行效率,原则上说,缓冲区越大,磁io的次数越少,执行速度就越快 默认mapreduceyarn要点:yarnyarn这样一来,yarnyarn上可以运行各程序,tez……yarn规范的资源请求机制即可Yarn就成为一个通用的资源调度平台,从此,企业中以前存在的各种运算集群都Partition就是对map输出的key进行分组,不同的组可以指定不同的reducePartition功能由partitioner的实现子类来实现Mapreduce的排 重MR中的常见排序机制:partial/total/secondaryMR排序是在map阶段输出之后,reduce(通过无reduce的MR程序示例观察)只针对key进行排序Key要实现 parable接口简单示例:对流量汇总数据进行倒序排partitalreducetask自动实现各输出文件total设置一个reducetask设置分区段partitioner,设置相应数量的reducetask,可以实现全局有序,job来统计数据分布规律,获取合适的区段划分,然后partitionerjob对整个数据集利用hadoop*自带的TotalOrderPartitioner**@author*publicclassTotalSortstaticclassTotalSortMapperextendsMapper<Text,Text,Text,Text>{OrderBeanbean=newOrderBean();protectedvoidmap(Textkey,Textvalue,Contextcontext)throwsIOException,InterruptedException{//Stringline=//String[]fields=//bean.set(fields[0],Double.parseDouble(fields[1]));context.write(key,value);}}protectedvoidreduce(Textkey,I ble<Text>values,Contextcontext)throwsIOException,InterruptedException{for(Textv:values){context.write(key,v);}}}publicstaticvoidmain(String[]args)throwsExceptionConfigurationconf=newConfiguration();Jobjob=Job.getInstance(conf); FileInputFormat.setInputPaths(job,newPath(args[0]));FileOutputFormat.setOutputPath(job,newPath(args[1]));RandomSamplerInputSampler.writePartitionFile(job,randomSampler);Configurationconf2job.getConfiguration();StringpartitionFile=}}secondarymapreducevalue考虑一个场景,需要取按key分组的最大value条目:通常,shuffle只是对key进行排序如果需要对value排序,则需要将value放到key中,但是此时,value就和原来keykeyreducerkey是一个一个到达reducerreducervalue的那一个,不好办,它会一个一个都输出去,除非自己弄一个缓存,将到达的组合key全部缓存起来然后只取第一个(或者弄一个标识?但是同一个reducerkeykey,无法判断标识)此时就可以用到secondarysort要有对组合key要有partitioner进行分区负载并行reducer parator来重定义valuelist聚合策略——这是关键其原理就是将相同key而不同组合key的数据进行聚合,从而把他们聚合成一组,然后在reducer中可以一次收到这一组key的组合key,并且,value最大的也就是在这一组中的第一个组合key会被选为迭代器valuelist的key,从而可以直接输出这个组合key,就实现了我们的需求示例:输出每个item定义一 用于控制shuffle过程中reduce端对kv@authorpublic parator paratorparator()super(OrderBean.class,}publicintparableparableb)OrderBeanabean=(OrderBean)a;OrderBeanbbean=(OrderBean)//将item_id相同的beanreturn}}定义订单信息订单信息bean,实现hadoop@authorpublicclassOrderBeanimplements privateTextitemid;privateDoubleWritablepublicOrderBean()}publicOrderBean(Textitemid,DoubleWritableamount){set(itemid,amount);}publicvoidset(Textitemid,DoubleWritableamount){this.itemid=itemid;this.amount=}publicTextgetItemid(){returnitemid;}publicDoubleWritablegetAmount(){returnamount;}publicintcompareTo(OrderBeano)intcmp= if(cmp==0){cmp }return}publicvoidwrite(DataOutputout)throwsIOException{}publicvoidreadFields(DataInputin)throwsIOException{StringreadUTF=in.readUTF();doublereadDouble=this.itemid=newthis.amount=new}}publicStringtoString()returnitemid.toString()+"\t"+}}自定义一个partitioner,以使相同id的bean发往相同reducepublicpublicclassItemIdPartitionerextendsPartitioner<OrderBean,publicintgetPartition(OrderBeankey,NullWritablevalue,intnumPartitions)//指定item_id相同的bean发往相同的reducer}} 定义mr利用secondarysort机制输出每种item@authorpublicclassSecondarySort OrderBean,NullWritable>{OrderBeanbean=newOrderBean();protectedvoidmap(LongWritablekey,Textvalue,Contextcontext)IOException,InterruptedExceptionStringline=String[]fields=StringUtils.split(line,context.write(bean,}}staticclassSecondarySortReducerextendsReducer<OrderBean,NullWritable,OrderBean,NullWritable>{ parator以后这里收到的kv数据就是:<1001 <100176.5>,null //此时,reducekeykvkv<1001//item

Contextcontext)throwsIOException,InterruptedException{context.write(key,NullWritable.get());}}publicstaticvoidmain(String[]args)throwsExceptionConfigurationconf=newConfiguration();Jobjob=Job.getInstance(conf);FileInputFormat.setInputPaths(job,newPath(args[0]));FileOutputFormat.setOutputPath(job,newPath(args[1]));//指定shuffle所使用 parator //指定shuffle所使用的partitioner}}shuffleShuffleshuffle是MRmaptask和reducetask3个操作:1、分区2、Sort根据key3、Combiner进行局部value整个shufflemaptaskcombiner分区/reducetask拉取mapreducevaluesreducekv(parator)中排序最前的kv的key传给reduce方法的入参mrmap一个inputsplit对应一个而inputsplit切片规划是由InputFormatInputSplits[] getSplits()方法,这个方法的逻辑可以自定义在默认情况下,由FileInputFormat来实现,它的逻辑(1)longlongminSize=Math.max(getFormatMinSplitSize(),getMinSplitSize(job));longmaxSize=getMaxSplitSize(job);}} //protectedlongcomputeSplitSize(longblockSize,longminSize,longmaxSize){returnMath.max(minSize,Math.min(maxSize,blockSize));}(2)构造切片信息对象,并放入InputSplitsblocksize1.1倍大小MapreducejoinReduceside通过将关联的条件作为map输出的keyjoin条件的数据并携带数据所来源的文件信息,发往同一个reducetask,在reduce中进行数据的串联publicpublicclassOrderJoin OrderJoinBean>{protectedvoidmap(LongWritablekey,Textvalue,Contextcontext)throwsIOException,InterruptedException{//Stringline=value.toString();String[]Stringline=value.toString();String[]fields=line.split("\t");//拿到Stringitemid=//获取到这一行所在的文件名(通过inpusplit)Stringname="你拿到的文件名";b,切分出3个字段)OrderJoinBeanbean=newOrderJoinBean();bean.set(null,null,null,null,null);context.write(newText(itemid),bean);}} OrderJoinBean,NullWritable>{protectedvoidreduce(Textkey,I ble<OrderJoinBean>beans,Contextcontext)throwsIOException,InterruptedException{//拿到的key是某一个itemid,//拿到的beans是来自于两类文件的 }}}注:也可利用二次排序的逻辑来实现reduceMapside可以将小表分发到所有的map节点,这样,map节点就可以在本地对自己所读到的大表数据进行join并输出最终结果可以大大提高join--mapper类中预先定义好小表,进行FileReaderin=null;BufferedReaderreader=HashMap<String,String>b_tab=newHashMap<String,String>();Stringlocalpath=null;Stringuirpath=protectedvoidsetup(Contextcontext)throwsIOException,InterruptedExceptionPath[]files=context.getLocalCacheFiles();localpath=files[0].toString();URI[]cacheFiles=in=newFileReader("b.txt");reader=newBufferedReader(in);Stringline=null;String[]fields=line.split(",");}}protectedvoidmap(LongWritablekey,Textvalue,Contextcontext)IOException,InterruptedException//maptask所负责的那一个切片数据(hdfs上String[]fields=Stringa_itemid=fields[0];Stringa_amount=fields[1];Stringb_name=//输出结果 98.9context.write(newText(a_itemid),newText(a_amount+"\t"+":"+localpath+"\t"+b_name));}}publicstaticvoidmain(String[]args)throwsException{Configurationconf=newConfiguration();Jobjob=Job.getInstance(conf);FileInputFormat.setInputPaths(job,newPath(args[0]));FileOutputFormat.setOutputPath(job,newPath(args[1]));//reducer //jartask节点的classpath}}mapreduceDistributed应用场景:mapside通过mapreduce框架将一个文件(本地/HDFS)分发到每一个运行时的task(maptask/reducetask)节点上(放到task进程所在的工 mapperreducer的代码内,直接使用本地文件JAVAAPI来这个文件首先在jobjob.addCacheFile(job.addCacheFile(newURI("hdfs://hadoop-然后在mapper或者reducerin=newFileReader("b.txt");readerin=newFileReader("b.txt");reader=newStringline=Mapreducemaptask数据切片与map任务数的机制isSplitable() 构造一个记录具体数据的逻辑是实现在LineRecordReader中(按行数据行起始偏移量作为key,value),比较特别的地方是:split多读一行(针对的是:非最末切片)InputFormatInputFormat--源码结构 在LineRecordReader中,对splitpublicpublicvoidinitialize(InputSplitTaskAttemptContextcontext)throwsIOException{FileSplitsplit=(FileSplit)genericSplit;Configurationjob=…//openthefileandseektothestartofthesplitfinalFileSystemfs=file.getFileSystem(job);fileInfileIn=CompressionCodeccodec=newCompressionCodecFactory(job).getCodec(file);if(null!=codec){………//我们总是将第一条记录抛弃(文件第一个split除外//因为我们总是在nextKeyValue()方法中跨split多读了一行(split除外if(start!=0)start+=in.readLine(newText(),0,}this.pos=}在LineRecordReader中,nextKeyValue方法总是跨splitpublicpublicbooleannextKeyValue()throwsIOException{if(key==null){key=new}if(value==null){value=newText();}intnewSize=//使用<=来 一while(getFilePosition()<=end||in.needAdditionalRecordAfterSplit()){newSize=in.readLine(value,maxLineLength,pos+=newSize;if(newSize<maxLineLength){….}它的切片逻辑跟TextInputformat完全不同:则将这些数据块形成一个切片,继承该过程,直到剩余数据块累加大小小于mapred.min.split.size.per.node,则将这些剩余数据块mapred.max.split.size,则将这些数据块形成一个切片,继承该过程,mapred.max.split.size,则进行下一步;aped.n.sp.sze.e.rckmapred.min.split.size.per.rack,则这些数据块留////fromaracknametothelistofblocksitHashMap<String,List<OneBlockInfo>>rackToBlocks=newHashMap<String,List<OneBlockInfo>>();//map fromablocktothenodesonwhichithasreplicasHashMap<OneBlockInfo,String[]>blockToNodes=new//map fromanodetothelistofblocksthatitcontainsHashMap<String,List<OneBlockInfo>>nodeToBlocks=newHashMap<String,List<OneBlockInfo>>();//populatealltheblocksforallfileslong//populatealltheblocksforallfileslongtotLength=0;for(inti=0;i<paths.length;i++){files[i]=newOneFileInfo(paths[i],totLength+=}}//ArrayList<OneBlockInfo>validBlocks=new//ArrayList<String>nodes=new//longcurSplitSize//processallnodesandcreatesplitsthatarelocaltoa//依次处理每个节点上的数据块for(I tor<Map.Entry<String,List<OneBlockInfo>>>iter=nodeToBlocks.entrySet().i tor();iter.hasNext();){Map.Entry<String,List<OneBlockInfo>>one=iter.next();List<OneBlockInfo>blocksInNode=//foreachblock,copyitintovalidBlocks.DeleteitfromblockToNodessothatthesameblockdoesnotappearin//twodifferentsplits.//blockToNodes出现在两个切片中forOneBlockInfooneblockblocksInNodeif(blockToNodes.containsKey(oneblock)){curSplitSize+=oneblock.length;//iftheaccumulatedsplitsizeexceedstheum,thencreatethissplit.//maxSizeif(maxSize!=&&curSplitSize>=maxSize) addCreatedSplit(job,splits,nodes,validBlocks);curSplitSize=0;}}}//iftherewereanyblocksleftoverandtheircombinedsize//largerthanminSplitNode,thencombinethemintoone//Otherwiseaddthembacktotheunprocessedpool.Itis//thattheywillbecombinedwithotherblocksfromthesameracklater////if(minSizeNode!=0&&curSplitSize>=minSizeNode){//createaninputsplitandaddittothesplitsarray addCreatedSplit(job,splits,nodes,validBlocks);}elsefor(OneBlockInfooneblock:validBlocks){for(OneBlockInfooneblock:validBlocks){}}curSplitSize=0;}//ifblocksinarackarebelowthespecifiedminimumsize,thenkeep//in'overflow'.Aftertheprocessingofallracksiscomplete,these//blockswillbecombinedinto//overflowBlocks用于保存“同一机架”过程处理之后剩余的数据块ArrayList<OneBlockInfo>overflowBlocksnewArrayList<OneBlockInfo>();ArrayList<String>racks=newArrayList<String>();//Processallracksoverandoveragainuntilthereisnomoreworktodo.while(blockToNodes.size()>0){//Createonesplitforthisrackbeforemovingovertothenext//Comebacktothisrackaftercreatingasinglesplitforeachof//remaining//Processoneracklocationatatime,Combineallpossibleblocks//resideonthisrackasonesplit.(constrainedbyminimum //split// teoverall//依次处理每个机架for(I tor<Map.Entry<String,List<OneBlockInfo>>>iter= tor();iter.hasNext();){Map.Entry<String,List<OneBlockInfo>>one=iter.next();List<OneBlockInfo>blocks=//foreachblock,copyitintovalidBlocks.Deleteitfrom//blockToNodessothatthesameblockdoesnotappearin//twodifferentsplits.booleancreatedSplit=false;//依次处理该机for(OneBlockInfooneblock:blocks){if(blockToNodes.containsKey(oneblock)){curSplitSize+=oneblock.length;//iftheaccumulatedsplitsizeexceedstheum,thencreatethissplit.//maxSizeifmaxSize&&curSplitSize>=maxSize) addCreatedSplit(job,splits,getHosts(racks),validBlocks);createdSplit=true;}}}//ifwecreatedasplit,thenjustgotothenextrackif(createdSplit){curSplitSize=0;}minSizeRack,则将这些数据块构成一个切片if(minSizeRack!=0&&curSplitSize>=minSizeRack){//ifthereisamimimumsizespecified,thencreateasingle addCreatedSplit(job,splits,getHosts(racks),validBlocks);}else//Therewereafewblocksinthisrackthatremainedtobe//Keepthemin'overflow'blocklist.Thesewillbecombined如果剩余数据块大小小于minSizeRack,则将这些数据块加入overflowBlocks}}curSplitSize=0;}}//Processalloverflowblocksfor(OneBlockInfooneblock:overflowBlocks){validBlocks.add(oneblock);curSplitSize+=oneblock.length;//Thismightcauseanexitingracklocationtobere-added,////Processalloverflowblocksfor(OneBlockInfooneblock:overflowBlocks){validBlocks.add(oneblock);curSplitSize+=oneblock.length;//Thismightcauseanexitingracklocationtobere-added,//butitshouldbeok.for(inti=0;i<oneblock.racks.length;i++){racks.add(oneblock.racks[i]);}//iftheaccumulatedsplitsizeexceeds um,//createthis//maxSizeifmaxSize0&&curSplitSize>=maxSize){//createaninputsplitandaddittothesplitsarray addCreatedSplit(job,//createaninputsplitandaddittothesplitsarray addCreatedSplit(job,splits,getHosts(racks), curSplitSize=0;}}//Process//Processanyremainingblocks,ifany.if(!validBlocks.isEmpty()){addCreatedSplit(job,splits,getHosts(racks),validBlocks);}CombineFileInputFormat形成切片过程中考虑数据本地性(同一节点、同一机架),首先处逐步减弱的。另外CombineFileInputFormat是抽象的,具体使用时需要自己实现getRecordReader方法。sequenceFile是hadoop中非常重要的一种数据格式sequenceFile文件内部的数据组织形式是:K-V对读入/写出为hadoop虽然FileInputFormat可以多个,但是有些场景下我们要处理的数据可能MultipleInputs,可以为不同的路径指定不同的mapper类来处理;SequenceFile格式文件staticstaticclassTextMapperAextendsMapper<LongWritable,Text,Text,LongWritable>{protectedvoidmap(LongWritablekey,Textvalue,Contextcontext)throwsIOException,InterruptedException{Stringline=value.toString();String[]words=line.split("");for(Stringw:words){}}}staticclassSequenceMapperBextendsMapper<Text,LongWritable,Text,LongWritable>{protectedvoidmap(Textkey,LongWritablevalue,Contextcontext) IOException,InterruptedException{}}jobMultipleInputsmapperpublicpublicstaticvoidmain(String[]args)throwsException{Configurationconf=newConfi

温馨提示

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

评论

0/150

提交评论