版权说明:本文档由用户提供并上传,收益归属内容提供方,若内容存在侵权,请进行举报或认领
文档简介
Hive内部表与外部表的区别?先来说下Hive中内部表与外部表的区别:
Hive创建内部表时,会将数据移动到数据仓库指向的路径;若创建外部表,仅记录数据所在的路径,
不对数据的位置做任何改变。在删除表的时候,内部表的元数据和数据会被一起删除,
而外部表只删除元数据,不删除数据。这样外部表相对来说更加安全些,数据组织也更加灵活,方便共享源数据。
需要注意的是传统数据库对表数据验证是schemaonwrite(写时模式),而Hive在load时是不检查数据是否
符合schema的,hive遵循的是schemaonread(读时模式),只有在读的时候hive才检查、解析具体的
数据字段、schema。
读时模式的优势是loaddata非常迅速,因为它不需要读取数据进行解析,仅仅进行文件的复制或者移动。
写时模式的优势是提升了查询性能,因为预先解析之后可以对列建立索引,并压缩,但这样也会花费要多的加载时间。
下面来看下Hive如何创建内部表:
1create
table
test(useridstring);2LOAD
DATAINPATH
'/tmp/result/20121213'
INTO
TABLE
testpartition(ptDate='20121213');这个很简单,不多说了,下面看下外部表:
01hadoopfs-ls/tmp/result/2012121402Found2items03-rw-r--r--
3junesupergroup
12402012-12-2617:15/tmp/result/20121214/part-0000004-rw-r--r--
1junesupergroup
12402012-12-2617:58/tmp/result/20121214/part-0000105--建表06create
EXTERNAL
table
IF
NOT
EXISTStest(useridstring)partitioned
by
(ptDatestring)ROWFORMATDELIMITEDFIELDSTERMINATED
BY
'\t';07--建立分区表,利用分区表的特性加载多个目录下的文件,并且分区字段可以作为where条件,更为重要的是08--这种加载数据的方式是不会移动数据文件的,这点和loaddata不同,后者会移动数据文件至数据仓库目录。09alter
table
test
add
partition(ptDate='20121214')location
'/tmp/result/20121214';--注意目录20121214最后不要画蛇添足加/*,我就是linuxshell用多了,加了这玩意,调试了一下午。。。注意:location后面跟的是目录,不是文件,hive会把整个目录下的文件都加载到表中:1create
EXTERNAL
table
IF
NOT
EXISTSuserInfo(id
int,sexstring,age
int,
name
string,emailstring,sdstring,edstring)
ROWFORMATDELIMITEDFIELDSTERMINATED
BY
'\t'
location
'/hive/dw';否则,会报错误:FAILED:Errorinmetadata:MetaException(message:Gotexception:org.apache.hadoop.ipc.RemoteExceptionjava.io.FileNotFoundException:Parentpathisnotadirectory:/hive/dw/record_2013-04-04.txt最后提下还有一种方式是建表的时候就指定外部表的数据源路径,但这样的坏处是只能加载一个数据源了:CREATEEXTERNALTABLEsunwg_test09(idINT,namestring)
ROWFORMATDELIMITED
FIELDSTERMINATEDBY‘\t’
LOCATION‘/sunwg/test08′;
上面的语句创建了一张名字为sunwg_test09的外表,该表有id和name两个字段,
字段的分割符为tab,文件的数据文件夹为/sunwg/test08
select*fromsunwg_test09;
可以查询到sunwg_test09中的数据。
在当前用户hive的根目录下找不到sunwg_test09文件夹。
此时hive将该表的数据文件信息保存到metadata数据库中。
mysql>select*fromTBLSwhereTBL_NAME=’sunwg_test09′;
可以看到该表的类型为EXTERNAL_TABLE。
mysql>select*fromSDSwhereSD_ID=TBL_ID;
在表SDS中记录了表sunwg_test09的数据文件路径为hdfs://hadoop00:9000/hjl/test08。
#hjl为hive的数据库名
实际上外表不光可以指定hdfs的目录,本地的目录也是可以的。
比如:
CREATEEXTERNALTABLEtest10(idINT,namestring)
ROWFORMATDELIMITED
FIELDSTERMINATEDBY‘\t’
2、Hbase的rowkey怎么创建比较好?列簇怎么创建比较好??IDCreateTimeNameCategoryUserID120120902中国好声音第1期综艺1220120904中国好声音第2期综艺1320120906中国好声音外卡赛综艺1420120908中国好声音第3期综艺1520120910中国好声音第4期综艺1620120912中国好声音选手采访综艺花絮2720120914中国好声音第5期综艺1820120916中国好声音录制花絮综艺花絮2920120918张玮独家专访花絮31020120920加多宝凉茶广告综艺广告4这里UserID应该对应另一张User表,暂不列出。我们只需知道UserID的含义:1代表浙江卫视;2代表好声音剧组;3代表XX微博;4代表赞助商。调用查询接口的时候将上述5个条件同时输入find(20120901,20121001,”中国好声音”,”综艺”,”浙江卫视”)。此时我们应该得到记录应该有第1、2、3、4、5、7条。第6条由于不属于“浙江卫视”应该不被选中。我们在设计RowKey时可以这样做:采用UserID+CreateTime+FileID组成RowKey,这样既能满足多条件查询,又能有很快的查询速度。需要注意以下几点:(1)每条记录的RowKey,每个字段都需要填充到相同长度。假如预期我们最多有10万量级的用户,则userID应该统一填充至6位,如000001,000002…(2)结尾添加全局唯一的FileID的用意也是使每个文件对应的记录全局唯一。避免当UserID与CreateTime相同时的两个不同文件记录相互覆盖。按照这种RowKey存储上述文件记录,在HBase表中是下面的结构:rowKey(userID6+time8+fileID6)namecategory….00000120120902000001000001201209040000020000012012090600000300000120120908000004000001201209100000050000012012091400000700000220120912000006000002201209160000080000032012091800000900000420120920000010怎样用这张表?在建立一个scan对象后,我们setStartRow(00000120120901),setEndRow(00000120120914)。这样,scan时只扫描userID=1的数据,且时间范围限定在这个指定的时间段内,满足了按用户以及按时间范围对结果的筛选。并且由于记录集中存储,性能很好。然后使用SingleColumnValueFilter(org.apache.hadoop.hbase.filter.SingleColumnValueFilter),共4个,分别约束name的上下限,与category的上下限。满足按同时按文件名以及分类名的前缀匹配。(注意:使用SingleColumnValueFilter会影响查询性能,在真正处理海量数据时会消耗很大的资源,且需要较长的时间)如果需要分页还可以再加一个PageFilter限制返回记录的个数。以上,我们完成了高性能的支持多条件查询的HBase表结构设计。用mapreduce怎么处理数据倾斜问题?map/reduce程序卡住的原因是什么?
2.根据原因,你是否能够想到更好的方法来解决?(企业很看重个人创作力)
map/reduce程序执行时,reduce节点大部分执行完毕,但是有一个或者几个reduce节点运行很慢,导致整个程序的处理时间很长,这是因为某一个key的条数比其他key多很多(有时是百倍或者千倍之多),这条key所在的reduce节点所处理的数据量比其他节点就大很多,从而导致某几个节点迟迟运行不完,此称之为数据倾斜。
用hadoop程序进行数据关联时,常碰到数据倾斜的情况,这里提供一种解决方法。
(1)设置一个hash份数N,用来对条数众多的key进行打散。
(2)对有多条重复key的那份数据进行处理:从1到N将数字加在key后面作为新key,如果需要和另一份数据关联的话,则要重写比较类和分发类(方法如上篇《hadoopjob解决大数据量关联的一种方法》)。如此实现多条key的平均分发。
intiNum=iNum%iHashNum;
StringstrKey=key+CTRLC+String.valueOf(iNum)+CTRLB+“B”;
(3)上一步之后,key被平均分散到很多不同的reduce节点。如果需要和其他数据关联,为了保证每个reduce节点上都有关联的key,对另一份单一key的数据进行处理:循环的从1到N将数字加在key后面作为新key
for(inti=0;i<iHashNum;++i){
StringstrKey=key+CTRLC+String.valueOf(i);
output.collect(newText(strKey),newText(strValues));}
以此解决数据倾斜的问题,经试验大大减少了程序的运行时间。但此方法会成倍的增加其中一份数据的数据量,以增加shuffle数据量为代价,所以使用此方法时,要多次试验,取一个最佳的hash份数值。
======================================
用上述的方法虽然可以解决数据倾斜,但是当关联的数据量巨大时,如果成倍的增长某份数据,会导致reduceshuffle的数据量变的巨大,得不偿失,从而无法解决运行时间慢的问题。
有一个新的办法可以解决成倍增长数据的缺陷:
在两份数据中找共同点,比如两份数据里除了关联的字段以外,还有另外相同含义的字段,如果这个字段在所有log中的重复率比较小,则可以用这个字段作为计算hash的值,如果是数字,可以用来模hash的份数,如果是字符可以用hashcode来模hash的份数(当然数字为了避免落到同一个reduce上的数据过多,也可以用hashcode),这样如果这个字段的值分布足够平均的话,就可以解决上述的问题。Hadoop框架如何优化?1.使用自定义Writable自带的Text很好用,但是字符串转换开销较大,故根据实际需要自定义Writable,注意作为Key时要实现WritableCompareable接口避免output.collect(newText(),newText())提倡key.set()value.set()output.collect(key,value)前者会产生大量的Text对象,使用完后Java垃圾回收器会花费大量的时间去收集这些对象
2.使用StringBuilder不要使用FormatterStringBuffer(
线程安全)StringBuffer尽量少使用多个append方法,适当使用+
3.使用DistributedCache加载文件比如配置文件,词典,共享文件,避免使用static变量
4.充分使用CombinerParttitionerComparator。Combiner:对map任务进行本地聚合Parttitioner:合适的Parttitioner避免reduce端负载不均Comparator:二次排序比如求每天的最大气温,map结果为日期:气温,若气温是降序的,直接取列表首元素即可
5.使用自定义InputFormat和OutputFormat
6.MR应避免静态变量:不能用于计数,应使用Counter大对象:MapList递归:避免递归深度过大超长正则表达式:消耗性能,要在map或reduce函数外编译正则表达式不要创建本地文件:变向的把HDFS里面的数据转移到TaskTracker,占用网络带宽不要大量创建目录和文件不要大量使用System.out.println,而使用Logger不要自定义过多的Counter,最好不要超过100个不要配置过大内存,mapred.child.java.opts-Xmx2000m是用来设置mapreduce任务使用的最大heap量7.关于map的数目map数目过大[创建和初始化map的开销],一般是由大量小文件造成的,或者dfs.block.size设置的太小,对于小文件可以archive文件或者Hadoopfs-merge合并成一个大文件.map数目过少,造成单个map任务执行时间过长,频繁推测执行,且容易内存溢出,并行性优势不能体现出来。dfs.block.size一般为256M-512M压缩的Text文件是不能被分割的,所以尽量使用SequenceFile,可以切分
8.关于reduce的数目reduce数目过大,产生大量的小文件,消耗大量不必要的资源,reduce数目过低呢,造成数据倾斜问题,且通常不能通过修改参数改变。可选方案:mapred.reduce.tasks设为-1变成AutoReduce。Key的分布,也在某种程度上决定了Reduce数目,所以要根据Key的特点设计相对应的Parttitioner避免数据倾斜
9.Map-side相关参数优化io.sort.mb(100MB):通常k个maptasks会对应一个buffer,buffer主要用来缓存map部分计算结果,并做一些预排序提高map性能,若map输出结果较大,可以调高这个参数,减少map任务进行spill任务个数,降低I/O的操作次数。若map任务的瓶颈在I/O的话,那么将会大大提高map性能。如何判断map任务的瓶颈?io.sort.spill.percent(0.8):spill操作就是当内存buffer超过一定阈值(这里通常是百分比)的时候,会将buffer中得数据写到Disk中。而不是等buffer满后在spill,否则会造成map的计算任务等待buffer的释放。一般来说,调整io.sort.mb而不是这个参数。io.sort.factor(10):map任务会产生很多的spill文件,而map任务在正常退出之前会将这些spill文件合并成一个文件,即merger过程,缺省是一次合并10个参数,调大io.sort.factor,减少merge的次数,减少DiskI/O操作,提高map性能。bine:通常为了减少map和reduce数据传输量,我们会制定一个combiner,将map结果进行本地聚集。这里combiner可能在merger之前,也可能在其之后。那么什么时候在其之前呢?当spill个数至少为bine指定的数目时同时程序指定了Combiner,Combiner会在其之前运行,减少写入到Disk的数据量,减少I/O次数。
10.压缩(时间换空间)MR中的数据无论是中间数据还是输入输出结果都是巨大的,若不使用压缩不仅浪费磁盘空间且会消耗大量网络带宽。同样在spill,merge(reduce也对有一个merge)亦可以使用压缩。若想在cpu时间和压缩比之间寻找一个平衡,LzoCodec比较适合。通常MR任务的瓶颈不在CPU而在于I/O,所以大部分的MR任务都适合使用压缩。
11.reduce-side相关参数优化reduce:copy->sort->reduce,也称shufflemapred.reduce.parellel.copies(5):任一个map任务可能包含一个或者多个reduce所需要数据,故一个map任务完成后,相应的reduce就会立即启动线程下载自己所需要的数据。调大这个参数比较适合map任务比较多且完成时间比较短的Job。mapred.reduce.copy.backoff:reduce端从map端下载数据也有可能由于网络故障,map端机器故障而失败。那么reduce下载线程肯定不会无限等待,当等待时间超过mapred.reduce.copy.backoff时,便放弃,尝试从其他地方下载。需注意:在网络情况比较差的环境,我们需要调大这个参数,避免reduce下载线程被误判为失败。io.sort.factor:recude将map结果下载到本地时,亦需要merge,如果reduce的瓶颈在于I/O,可尝试调高增加merge的并发吞吐,提高reduce性能、mapred.job.shuffle.input.buffer.percent(0.7):reduce从map下载的数据不会立刻就写到Disk中,而是先缓存在内存中,mapred.job.shuffle.input.buffer.percent指定内存的多少比例用于缓存数据,内存大小可通过mapred.child.java.opts来设置。和map类似,buffer不是等到写满才往磁盘中写,也是到达阈值就写,阈值由mapred.job,shuffle.merge.percent来指定。若Reduce下载速度很快,容易内存溢出,适当增大这个参数对增加reduce性能有些帮助。mapred.job.reduce.input.buffer.percent(0):当Reduce下载map数据完成之后,就会开始真正的reduce的计算,reduce的计算必然也是要消耗内存的,那么在读物reduce所需要的数据时,同样需要内存作为buffer,这个参数是决定多少的内存百分比作为buffer。默认为0,也就是说reduce全部从磁盘读数据。若redcue计算任务消耗内存很小,那么可以设置这个参数大于0,使一部分内存用来缓存数据。Hbase内部是什么机制?深入分析HBaseRPC(Protobuf)实现机制\o"Binospace"Binospace
2013-08-02
2730
阅读背景在HMaster、RegionServer内部,创建了RpcServer实例,并与Client三者之间实现了Rpc调用,HBase0.95内部引入了Google-Protobuf作为中间数据组织方式,并在Protobuf提供的Rpc接口之上,实现了基于服务的Rpc实现,本文详细阐述了HBase-Rpc实现细节。HBase的RPCProtocol
在HMaster、RegionServer内部,实现了rpc多个protocol来完成管理和应用逻辑,具体如下protocol如下:HMaster支持的Rpc协议:
MasterMonitorProtocol,Client与Master之间的通信,Master是RpcServer端,主要实现HBase集群监控的目的。MasterAdminProtocol,Client与Master之间的通信,Master是RpcServer端,主要实现HBase表格的管理。例如TableSchema的更改,Table-Region的迁移、合并、下线(Offline)、上线(Online)以及负载平衡,以及Table的删除、快照等相关功能。RegionServerStatusProtoco,RegionServer与Master之间的通信,Master是RpcServer端,负责提供RegionServer向HMaster状态汇报的服务。RegionServer支持的Rpc协议:ClientProtocol,Client与RegionServer之间的通信,RegionServer是RpcServer端,主要实现用户的读写请求。例如get、multiGet、mutate、scan、bulkLoadHFile、执行Coprocessor等。AdminProtocols,Client与RegionServer之间的通信,RegionServer是RpcServer端,主要实现Region、服务、文件的管理。例如storefile信息、Region的操作、WAL操作、Server的开关等。(备注:以上提到的Client可以是用户Api、也可以是RegionServer或者HMaster)
HBase-RPC实现机制分析RpcServer配置三个队列:1)普通队列callQueue,绝大部分Call请求存在该队列中:callQueue上maxQueueLength为${ipc.server.max.callqueue.length},默认是${hbase.master.handler.count}*DEFAULT_MAX_CALLQUEUE_LENGTH_PER_HANDLER,目前0.95.1中,每个Handler上CallQueue的最大个数默认值(DEFAULT_MAX_CALLQUEUE_LENGTH_PER_HANDLER)为10。2)优先级队列:PriorityQueue。如果设置priorityHandlerCount的个数,会创建与callQueue相当容量的queue存储Call,该优先级队列对应的Handler的个数由rpcServer实例化时传入。3)拷贝队列:replicationQueue。由于RpcServer由HMaster和RegionServer共用,该功能仅为RegionServer提供,queue的大小为${ipc.server.max.callqueue.size}指定,默认为1024*1024*1024,handler的个数为hbase.regionserver.replication.handler.count。RpcServer由三个模块组成:Listener===Queue===Responder
这里以HBaseAdmin.listTables为例,分析一个Rpc请求的函数调用过程:1)RpcClient创建一个BlockingRpcChannel。2)以channel为参数创建执行RPC请求需要的stub,此时的stub已经被封装在具体Service下,stub下定义了可执行的rpc接口。3)stub调用对应的接口,实际内部channel调用callBlockingMethod方法。RpcClient内实现了protobuf提供的BlockingRpcChannel接口方法callBlockingMethod,
@OverridepublicMessagecallBlockingMethod(MethodDescriptormd,RpcControllercontroller,Messageparam,MessagereturnType)throwsServiceException{returnthis.rpcClient.callBlockingMethod(md,controller,param,returnType,this.ticket,this.isa,this.rpcTimeout);}通过以上的实现细节,最终转换成rpcClient的调用,使用MethodDescriptor封装了不同rpc函数,使用Message基类可以接收基于Message的不同的Request和Response对象。4)RpcClient创建Call对象,查找或者创建合适的Connection,并唤醒Connection。5)Connection等待Call的Response,同时rpcClient调用函数中,会使用connection.writeRequest(Callcall)将请求写入到RpcServer网络流中。6)等待Call的Response,然后层层返回给更上层接口,从而完成此次RPC调用。RPCServer收到的Rpc报文的内部组织如下:Magic(4Byte)Version(1Byte)AuthMethod(1Byte)ConnectionHeaderLength(4Byte)ConnectionHeaderRequest“HBas”验证RpcServer的CURRENT_VERSION与RPC报文一致目前支持三类:AuthMethod.SIMPLEAuthMethod.KERBEROSAuthMethod.DIGESTRPC.proto定义
RPCProtos.ConnectionHeader
messageConnectionHeader{
optionalUserInformationuserInfo=1;
optionalstringserviceName=2;
//Cellblockcodecwewillusesendingoveroptionalcellblocks.
Serverthrowsexception
//ifcannotdeal.
optionalstringcellBlockCodecClass=3[default="org.apache.hadoop.hbase.codec.KeyValueCodec"];
//Compressorwewilluseifcellblockiscompressed.
Serverwillthrowexceptionifnotsupported.
//Classmustimplementhadoop’sCompressionCodecInterface
optionalstringcellBlockCompressorClass=4;
}
序列化之后的数据整个Request存储是经过编码之后的byte数组,包括如下几个部分:RequestHeaderLength(RawVarint32)RequestHeaderParamSize(RawVarint32)ParamCellScannerRPC.proto定义:
messageRequestHeader{
//MonotonicallyincreasingcallIdtokeeptrackofRPCrequestsandtheirresponse
optionaluint32callId=1;
optionalRPCTInfotraceInfo=2;
optionalstringmethodName=3;
//Iftrue,thenapbMessageparamfollows.
optionalboolrequestParam=4;
//Ifpresent,thenanencodeddatablockfollows.
optionalCellBlockMetacellBlockMeta=5;
//TODO:Haveclientspecifypriority
}
序列化之后的数据
并从Header中确认是否存在Param和CellScanner,如果确认存在的情况下,会继续访问。Protobuf的基本类型Message,
Request的Param继承了Message,
这个需要获取的Method类型决定。从功能上讲,RpcServer上包含了三个模块,1)Listener。包含了多个Reader线程,通过Selector获取ServerSocketChannel接收来自RpcClient发送来的Connection,并从中重构Call实例,添加到CallQueue队列中。
”IPCServerlisteneron60021″daemonprio=10tid=0x00007f7210a97800nid=0x14c6runnable[0x00007f720e8d0000]
java.lang.Thread.State:RUNNABLE
atsun.nio.ch.EPollArrayWrapper.epollWait(NativeMethod)
atsun.nio.ch.EPollArrayWrapper.poll(EPollArrayWrapper.java:210)
atsun.nio.ch.EPollSelectorImpl.doSelect(EPollSelectorImpl.java:65)
atsun.nio.ch.SelectorImpl.lockAndDoSelect(SelectorImpl.java:69)
-locked<0x00000000c43cae68>(asun.nio.ch.Util$2)
-locked<0x00000000c43cae50>(ajava.util.Collections$UnmodifiableSet)
-locked<0x00000000c4322ca8>(asun.nio.ch.EPollSelectorImpl)
atsun.nio.ch.SelectorImpl.select(SelectorImpl.java:80)
atsun.nio.ch.SelectorImpl.select(SelectorImpl.java:84)
atorg.apache.hadoop.hbase.ipc.RpcServer$Listener.run(RpcServer.java:646)2)Handler。负责执行Call,调用Service的方法,然后返回Pair<Message,CellScanner>“IPCServerhandler0on60021″daemonprio=10tid=0x00007f7210eab000nid=0x14c7waitingoncondition[0x00007f720e7cf000]
java.lang.Thread.State:WAITING(parking)
atsun.misc.Unsafe.park(NativeMethod)
-parkingtowaitfor
<0x00000000c43cad90>(ajava.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject)
atjava.util.concurrent.locks.LockSupport.park(LockSupport.java:156)
atjava.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject.await(AbstractQueuedSynchronizer.java:1987)
atjava.util.concurrent.LinkedBlockingQueue.take(LinkedBlockingQueue.java:399)
atorg.apache.hadoop.hbase.ipc.RpcServer$Handler.run(RpcServer.java:1804)3)Responder。负责把Call的结果返回给RpcClient。
”IPCServerResponder”daemonprio=10tid=0x00007f7210a97000nid=0x14c5runnable[0x00007f720e9d1000]
java.lang.Thread.State:RUNNABLE
atsun.nio.ch.EPollArrayWrapper.epollWait(NativeMethod)
atsun.nio.ch.EPollArrayWrapper.poll(EPollArrayWrapper.java:210)
atsun.nio.ch.EPollSelectorImpl.doSelect(EPollSelectorImpl.java:65)
atsun.nio.ch.SelectorImpl.lockAndDoSelect(SelectorImpl.java:69)
-locked<0x00000000c4407078>(asun.nio.ch.Util$2)
-locked<0x00000000c4407060>(ajava.util.Collections$UnmodifiableSet)
-locked<0x00000000c4345b68>(asun.nio.ch.EPollSelectorImpl)
atsun.nio.ch.SelectorImpl.select(SelectorImpl.java:80)
atorg.apache.hadoop.hbase.ipc.RpcServer$Responder.doRunLoop(RpcServer.java:833)
atorg.apache.hadoop.hbase.ipc.RpcServer$Responder.run(RpcServer.java:816)RpcClient为Rpc请求建立Connection,通过Connection将Call发送RpcServer,然后RpcClient等待结果的返回。
思考1)为什么HBase新版本使用了Protobuf,并实现RPC接口?HBase是Hadoop生态系统内重要的分布式数据库,Hadoop2.0广泛采用Protobuf作为中间数据组织方式,整个系统内Wire-Compatible的统一需求。2)HBase内部实现的Rpc框架对于服务性能的影响?目前使用Protobuf作为用户请求和内部数据交换的数据格式,采用更为紧缩编码格式,能够提高传输数据的效率。但是,有些优化仍然可以在该框架内探索:实现多个Request复用Connection(把多个短连接合并成一个长连接);在RpcServer内创建多个CallQueue,分别处理不同的Service,分离管理逻辑与应用逻辑的队列,保证互不干扰;Responder单线程的模式,是否高并发应用的瓶颈所在?是否可以分离Read/Write请求占用的队列,以及处理的handler,从而使得读写性能能够更加平衡?针对读写应用的特点,在RpcServer层次内对应用进行分级,建立不同优先级的CallQueue,按照Hadoop-FairScheduler的模式,然后配置中心调度(类似OMega或者Spallow轻量化调度方案),保证实时应用的低延迟和非实时应用的高吞吐。优先级更好的Call会优先被调度给Handler,而非实时应用可以实现多个Call的合并操作,从而提高吞吐。3)Protobuf内置编码与传统压缩技术是否可以配合使用?使用tcpdump获取了一段HMaster得到的RegionServer上报来的信息:以上的信息几乎是明文出现在tcp-ip连接中,因此,是否在Protobuf-RPC数据格式采取一定的压缩策略,会给scan、multiGet等数据交互较为密集的应用提供一种优化的思路。我们在开发分布式计算job的,是否可以去掉reduce()阶段?hdfs的数据压缩算法在海量存储系统中,好的压缩算法可以有效的降低存储开销,减轻运营成本。Hadoop基本上已经是主流的存储系统了,但是由于本身是Java实现,在压缩性能上受到语言的限制;此外该压缩算法还得支持HadoopMapreduce的split机制,不然压缩后文件只能被一个mapper处理,会大大的影响效率。Twitter之前实现过一个lzo的压缩框架,很好的将lzo引入到hadoop中。借助其原理我实现了一个更加通用的压缩框架,姑且叫做Nsm(NativeSplittableMulti-method)Compression,可以方便的将一个Native(C/C++)实现的Encoder/Decoder整合到HadoopCompressionIO机制中,并支持Mapreduce时的切分;此外我还基于lzma2实现了一个更高效的通用日志压缩算法,压缩比为zlib的1/2,压缩速度为原lzma2的两倍。
首先看最基本的压缩算法。Bigtable里提到了一种针对网页的longcommonstring压缩,在网页按url聚集后可以达到十几倍的压缩比。然而在日志系统中,由于日志本身有过优化,很难出现网页那样大段重复的情况,也无法进行相似日志的聚集,longcommonstring的优势体现不出来,反而不如zlib这样的短窗口压缩。另一方面,基于可读性的考虑,日志中整数,md5,timestamp,ip这样的数据往往以文本形式表示。因此,一个自然而然的想法是先把文件预处理,将上述数据由文本转换为二进制,再交给通用压缩算法进行压缩。有意思的是,经过预处理后的原文件虽然往往大小可以缩小一半,但再由通用压缩算法压缩后的结果,却和直接压缩的大小相差无几;这主要是通用压缩算法都会对结果进行哈夫曼或者算术编码,因此预处理的作用并不明显(往往只能降低10%左右)。
虽然预处理对最终压缩比影响不大,但是由于预处理速度快,并减少了通用压缩算法要处理的数据量,因此往往可以提高压缩速度,特别是对lzma2这种慢速压缩算法。在我的实现里,预处理平均可以将原文件大小缩小到1/2,预处理+lzma2对比单纯使用lzma2,可以将压缩速度从2M/s提高到4M/s,压缩比由19%提高到17%;当然解压缩速度也从100M/s降低到50M/s,这也是个代价。预处理还有一个好处,就是可以把日志中不同类型的数据聚集在一起(类似按列存储的数据库),虽然我还没有实现,但相信会对压缩比有很大提升。(注:PPMd和Lzma2是我以为最好的通用文本压缩算法,不过预处理和PPMd结合的并不好,我猜测是因为PPMd是纯粹的算术编码,预处理分散了概率分布,起到了副作用;此外PPMd的解压缩度只有10M/s,也不可接受)。
再看Hadoop的CompressionIO机制,用户要实现一个自定义的Codec,用来创建Compressor,Decompressor,InputStream(DecompressStream)和OutputStream(CompressStream);Hadoop使用该InputStream和OutputStream来读写文件。为了支持Mapreduce时的split,需要实现blockbased的InputStream/OutputStream,即以block为单位进行数据压缩,并且还能让每个split都恰好从block头开始,才能让解压缩器识别。所以,可以用额外的Indexer程序为每个压缩文件生成一个index文件,记录每个block的offset;然后再实现一个自定义的InputFormat来实现切分功能,切分逻辑很简单,就是读取文件对应的index,把父类方法完成的splits(通常是按chunk64M划分)对齐到block的开始;对于没有index的压缩文件,则只能以其整体作为一个split。
最后就来看框架本身的实现了。首先在Native实现中,定义了Encoder/Decoder两个接口,为了简单,Encoder每次都会把传入的数据独立encode成一个block,Decoder也只能接受一个完整block的数据。EncodeUtil管理所有的encoder,用户将原始数据和压缩方法ID传给Util,Util再找到对应的Encoder进行数据压缩,并添加一个blockheader;该header包括magicnumber,raw/encodedlength,raw/encodedchecksum和压缩算法ID。同样,也有一个对应的DecodeUtil管理所有的decoder,util首先解析blockheader,得到算法ID并验证checksum后,调用对应的Decoder进行解压。所以,想要添加一个新算法,只需要实现Encoder/Decoder并在Util中注册即可。
而Hadoop端的实现也很简单,只用实现之前讲过的Codec,Compressor,Decompressor,InputStream,OutputStream,Indexer,InputFormat即可。布署的时候,添加Codec到hadoop-site.xml(core-site.xml),如下
pression.codecs
press.GzipCodec,press.DefaultCodec,press.BZip2Codec,pression.nsm.NsmCodec
pression.codec.nsm.class
pression.nsm.NsmCodec以及在hadoop-site.xml(mapred-site.xml)中添加
mapred.child.env
JAVA_LIBRARY_PATH=/path/to/your/hadoop/lib/native此外,还需要将相关的nativelib(libnsm.solibp7z.so)添加到Hadoop的lib/native中,并修改hadoop-config.sh,将LD_LIBRARY_PATH=/path/to/your/hadoop/lib/nativeexport即可(需要重启hadoop)。---------------另:预处理的实现和改进由于日志本身不一定是格式化对齐的,即使对齐,一个field中也可能含有多个可转换的数据,所以这其实是一个模式识别-转换的问题;另一方面,日志本身又有一些通用的格式可以利用。在我的实现里,是预先定义好一些splittoken,并给0~255每个byte归于一种类型;在扫描过程中,每遇到一个token,就检查该段数据类型并进行相应的转换。这样把基于byte的状态机变成基于segment的状态机,实现上更加简单高效。也正是因此,在该基础上把同一类型的segment存储在一起也很容易实现,只不过如何做到高效(时间,空间)想必在工程上也需要不小的功夫。hadoop中压缩知识点总结
(转载)
(2012-11-1309:17:06)转载▼标签:
杂谈Hadoop学习笔记(2)
———数据压缩问题
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
在hadoop中文件的压缩带来了两大好处:
(1)它减少了存储文件所需的空间;(2)加快了数据在网络上或者从磁盘上或到磁盘上的传输速度;
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
所有的压缩算法都显示出一种时间空间的权衡:更快的压缩和解压速度通常会耗费更多的空间;~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
编码/解码器用以执行压缩解压算法。在Hadoop里,编码/解码器是通过一个压缩解码器接口实现的;DEFLATE
press.DefaultCodecgzip
press.GzipCodecbzip2
press.BZip2CodecLZO
pression.lzo.LzopCodec
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
CompressionCodec对流进行压缩和解压缩CompressionCodec有两个方法可以用于轻松压缩或解压缩数据:如果想对一个正在被写入的输出流的数据进行压缩,我们可以使用createOutStream(OutputStreamout)方法创建一个CompressionOutputStream,将其压缩格式写入底层的流;反之,要想对从输入流读取而来的数据进行解压缩,则调用createOutStream(InoutStreamin)方法,从而获得一个compressionInputStream,从而获得一个CompressionInputStream,从而从底层的流读取未压缩的数据;
CompressionInputStream和CompressionOutputStream类似于java.util.zip.DeflaterOutStream和java.util.zip.DeflaterOutStream,前两者还可以提供重置其底层压缩和解压缩功能。
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
用CompressionCodecFactory方法来推断CompressionCodecCompressionCodeFactory提供了getCodec()方法,从而将文件扩展名映射到相应的CompressionCodec;从方法接受一个Path对象;(对这种东西必须拿出具体代码来说,以后对实例分析的时候会重新提出的!)~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
一个压缩的程序:publicclassFileDecompressor{publicstaticvoidmain(String[]args)throwsException{Stringuri=args[0];Configurationconf=newConfiguration();FileSystemfs=FileSystem.get(URI.create(uri),conf);PathinputPath=newPath(uri);CompressionCodecFactoryfactory=newCompressionCodecFactory(conf);
//检测CompressionCodeccodec=factory.getCodec(inputPath);if(codec==null){System.err.println("Nocodecfoundfor"+uri);System.exit(1);}StringoutputUri=CompressionCodecFactory.removeSuffix(uri,codec.getDefaultExtension());//移除后缀名,恢复InputStreamin=null;OutputStreamout=null;try{
in=codec.createInputStream(fs.open(inputPath));//CompressionCodec对流进行压缩out=fs.create(newPath(outputUri));IOUtils.copyBytes(in,out,conf);}finally{IOUtils.closeStream(in);IOUtils.closeStream(out);}}~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~还有两个问题没有解决:(1)本地库的压缩解码;(这个是不懂)(2)压缩和输入分割问题:(这个是必须结合mapreduce实例来分解的,现在就放放呗~)
其实我们在编写mapreduce程序如果数据过大(闲得没事想用用压缩)需要涉及压缩数据事,我们可以重写和调用压缩和解压缩class来实现需要的过程mapreduce的调度模式一个MapReduce作业的生命周期大体分为5个阶段
【1】
:1.
作业提交与初始化2.
任务调度与监控3.任务运行环境准备4.任务执行5.作业完成我们假设JobTracker已经启动,那么调度器是怎么启动的?JobTracker在启动时有以下代码:JobTrackertracker=startTracker(newJobConf());tracker.offerService();其中offerService方法负责启动JobTracker提供的各个服务,有这样一行代码:taskScheduler.start();taskScheduler即为任务调度器。start方法是抽象类TaskScheduler提供的接口,用于启动调度器。每个调度器类都要继承TaskScheduler类。回忆一下,调度器启动时会将各个监听器对象注册到JobTracker,以FIFO调度器JobQueueTaskScheduler为例:@Overridepublicsynchronizedvoidstart()throwsIOException{super.start();taskTrackerManager.addJobInProgressListener(jobQueueJobInProgressListener);eagerTaskInitializationListener.setTaskTrackerManager(taskTrackerManager);eagerTaskInitializationListener.start();taskTrackerManager.addJobInProgressListener(eagerTaskInitializationListener);}这里注册了两个监听器,其中eagerTaskInitializationListener负责作业初始化,而jobQueueJobInProgressListener则负责作业的执行和监控。当有作业提交到JobTracker时,JobTracker会执行所有订阅它消息的监听器的jobAdded方法。对于eagerTaskInitializationListener来说:
@OverridepublicvoidjobAdded(JobInProgressjob){synchronized(jobInitQueue){jobInitQueue.add(job);resortInitQueue();jobInitQueue.notifyAll();}}提交的作业的JobInProgress对象被添加到作业初始化队列jobInitQueue中,并唤醒初始化线程(若原来没有作业可以初始化):classJobInitManagerimplementsRunnable{publicvoidrun(){JobInProgressjob=null;while(true){try{synchronized(jobInitQueue){while(jobInitQueue.isEmpty()){jobInitQueue.wait();}job=jobInitQueue.remove(0);}threadPool.execute(newInitJob(job));}catch(InterruptedExceptiont){LOG.info("JobInitManagerThreadinterrupted.");break;}}threadPool.shutdownNow();}}这种工作方式是一种“生产者-消费者”模式:作业初始化线程是消费者,而监听器eagerTaskInitializationListener是生产者。这里可以有多个消费者线程,放到一个固定资源的线程池中,线程个数通过mapred.jobinit.threads参数配置,默认为4个。下面我们重点来看调度器中的另一个监听器。
jobQueueJobInProgressListener对象在调度器中初始化时连续执行了两个构造器完成初始化:publicJobQueueJobInProgressListener(){this(newTreeMap<JobSchedulingInfo,JobInProgress>(FIFO_JOB_QUEUE_COMPARATOR));}/***Forclientsthatwanttoprovidetheirownjobpriorities.*@paramjobQueueAcollectionwhoseiteratorreturnsjobsinpriorityorder.*/protectedJobQueueJobInProgressListener(Map<JobSchedulingInfo,JobInProgress>jobQueue){this.jobQueue=Collections.synchronizedMap(jobQueue);}其中,第一个构造器调用重载的第二个构造器。可以看到,调度器使用一个队列jobQueue来保存提交的作业。这个队列使用一个TreeMap对象实现,TreeMap的特点是底层使用红黑树实现,可以按照键来排序,并且由于是平衡树,效率较高。作为键的是一个JobSchedulingInfo对象,作为值就是提交的作业对应的JobInProgress对象。另外,由于TreeMap本身不是线程安全的,这里使用了集合类的同步方法构造了一个线程安全的Map。使用带有排序功能的数据结构的目的是使作业在队列中按照优先级的大小排列,这样每次调度器只需从队列头部获得作业即可。作业的顺序由优先级决定,而优先级信息包含在JobSchedulingInfo对象中:staticclassJobSchedulingInfo{privateJobPrioritypriority;privatelongstartTime;privateJobIDid;...}该对象包含了作业的优先级、ID和开始时间等信息。在Hadoop中,作业的优先级有以下五种:VERY_HIGH、HIGH、NORMAL、LOW、VERY_LOW。这些字段是通过作业的JobStatus对象初始化的。由于该对象作为TreeMap的键,因此要实现自己的equals方法和hashCode方法:@Overridepublicbooleanequals(Objectobj){if(obj==null||obj.getClass()!=JobSchedulingInfo.class){returnfalse;}elseif(obj==this){returntrue;}elseif(objinstanceofJobSchedulingInfo){JobSchedulingInfothat=(JobSchedulingInfo)obj;return(this.id.equals(that.id)&&this.startTime==that.startTime&&this.priority==that.priority);}returnfalse;}我们看到,两个JobSchedulingInfo对象相等的条件是类型一致,并且作业ID、开始时间和优先级都相等。hashCode的计算比较简单:@OverridepublicinthashCode(){return(int)(id.hashCode()*priority.hashCode()+startTime);}注意,监听器的第一个构造器有一个比较器参数,用于定义
JobSchedulingInfo的比较方式:staticfinalComparator<JobSchedulingInfo>FIFO_JOB_QUEUE_COMPARATOR=newComparator<JobSchedulingInfo>(){publicintcompare(JobSchedulingInfoo1,JobSchedulingInfoo2){intres=o1.getPriority().compareTo(o2.getPriority());if(res==0){if(o1.getStartTime()<o2.getStartTime()){res=-1;}else{res=(o1.getStartTime()==o2.getStartTime()?0:1);}}if(res==0){res=o1.getJobID().compareTo(o2.getJobID());}returnres;}};从上面看出,首先比较作业的优先级,若优先级相等则比较开始时间(FIFO),若再相等则比较作业ID。
我们在实现自己的调度器时可能要定义自己的作业队列,那么作业在队列中的顺序(即
JobSchedulingInfo的比较器
)就要仔细定义,这是调度器能够正常运行基础。Hadoop中的作业调度采用pull方式,即TaskTracker定时向JobTracker发送心跳信息索取一个新的任务,这些信息包括数据结点上作业和任务的运行情况,以及该TaskTracker上的资源使用情况。JobTracker会依据以上信息更新作业队列的状态,并调用调度器选择一个或多个任务以心跳响应的形式返回给TaskTracker。从上面描述可以看出,JobTracker和taskScheduler之间的互相利用关系:前者利用后者为TaskTracker分配任务;后者利用前者更新队列和作业信息。接下来,我们一步步详述该过程。首先,当一个心跳到达JobTracker时(实际上这是一个来自TaskTracker的远程过程调用
heartbeat方法
,协议接口是InterTrackerProtocol),会执行两种动作:更新状态和下达命令
【1】
。下达命令稍后关注。有关更新状态的一些代码片段如下:if(!processHeartbeat(status,initialContact,now)){if(prevHeartbeatResponse!=null){trackerToHeartbeatResponseMap.remove(trackerName);}returnnewHeartbeatResponse(newResponseId,newTaskTrackerAction[]{newReinitTrackerAction()});}具体的心跳处理,由私有函数processHeartbeat完成。该函数中有以下两个方法调用:updateTaskStatuses(trackerStatus);updateNodeHealthStatus(trackerStatus,timeStamp);分别用来更新任务的状态和结点的健康状态。在第一个方法中有下面代码片段:TaskInProgresstip=taskidToTIPMap.get(taskId);//Checkifthetipisknowntothejobtracker.Incaseofarestarted//jt,sometasksmightjoininlaterif(tip!=null||hasRestarted()){if(tip==null){tip=job.getTaskInProgress(taskId.getTaskID());job.addRunningTaskToTIP(tip,taskId,status,false);}//UpdatethejobandinformthelistenersifnecessaryJobStatusprevStatus=(JobStatus)job.getStatus().clone();//CloneTaskStatusobjecthere,becauseJobInProgress//orTaskInProgresscanmodifythisobjectand//thechangesshouldnotgetreflectedinTaskTrackerStatus.//AnoldTaskTrackerStatusisusedlaterincountMapTasks,etc.job.updateTaskStatus(tip,(TaskStatus)report.clone());JobStatusnewStatus=(JobStatus)job.getStatus().clone();//Updatethelistenersifanincompletejobcompletesif(prevStatus.getRunState()!=newStatus.getRunState()){JobStatusChangeEventevent=newJobStatusChangeEvent(job,EventType.RUN_STATE_CHANGED,
温馨提示
- 1. 本站所有资源如无特殊说明,都需要本地电脑安装OFFICE2007和PDF阅读器。图纸软件为CAD,CAXA,PROE,UG,SolidWorks等.压缩文件请下载最新的WinRAR软件解压。
- 2. 本站的文档不包含任何第三方提供的附件图纸等,如果需要附件,请联系上传者。文件的所有权益归上传用户所有。
- 3. 本站RAR压缩包中若带图纸,网页内容里面会有图纸预览,若没有图纸预览就没有图纸。
- 4. 未经权益所有人同意不得将文件中的内容挪作商业或盈利用途。
- 5. 人人文库网仅提供信息存储空间,仅对用户上传内容的表现方式做保护处理,对用户上传分享的文档内容本身不做任何修改或编辑,并不能对任何下载内容负责。
- 6. 下载文件中如有侵权或不适当内容,请与我们联系,我们立即纠正。
- 7. 本站不保证下载资源的准确性、安全性和完整性, 同时也不承担用户因使用这些下载资源对自己和他人造成任何形式的伤害或损失。
最新文档
- 采煤课程设计实习报告
- 车间调度算法设计课程设计
- 健康手环蓝牙BLE方案开发课程设计
- PCA降维生物信息课程设计
- 基于NLP的情感分析工具实现方法课程设计
- 深度强化学习游戏AI智能体设计课程设计
- ug应用课程设计小结
- 基于PID的直流电机效率控制课程设计
- 模拟 IC 设计工程师考试试卷及答案
- 彩笔画橘子课程设计
- 2026年秋季学期新版人教版小学数学五年级上册教学计划附进度表
- CSCO胰腺癌诊疗指南(2026版)完整核心要点
- 《2023CSCO胃癌诊疗指南》解读
- 2026贵州省农业发展集团有限责任公司招录(第一批)岗位65人备考题库及完整答案详解
- 权利的正当性:理论、基础与实践探究
- 传媒行业内容审核标准(标准版)
- UG练习图纸大全-65张-绝对受用
- 安检金属探测器调试工程师岗位招聘考试试卷及答案
- ASME B16.10-2022 阀门结构长度(中英文参考版)
- 吊具管理制度规范
- 应用大地测量学 课件全套 第1-8章 绪论-空间大地测量
评论
0/150
提交评论