2015年7月15日星期三

hive优化使用技巧

参考:http://itindex.net/detail/53619-hive-优化

Hive是将符合SQL语法的字符串解析生成可以在Hadoop上执行的MapReduce的工具。使用Hive尽量按照分布式计算的一些特点来设计sql,和传统关系型数据库有区别,所以需要去掉原有关系型数据库下开发的一些固有思维。

基本原则:

1:尽量尽早地过滤数据,减少每个阶段的数据量,对于分区表要加分区,同时只选择需要使用到的字段。
join时,对于where条件,可以使用子查询设计join。
例如:
select .. from (select .. from ..) A join (select .. from ..) B on ..

2:尽量原子化操作,尽量避免一个SQL包含复杂逻辑,可以使用中间表来完成复杂的逻辑.
3:单个SQL所起的JOB个数尽量控制在5个以下.

4:慎重使用mapjoin,一般行数小于2000行,大小小于1M(扩容后可以适当放大)的表才能使用,小表要注意放在join的左边,否则会引起磁盘和内存的大量消耗.

5:写SQL要先了解数据本身的特点,如果有join ,group操作的话,要注意是否会有数据倾斜,如果出现数据倾斜,应当做如下处理:

set hive.exec.reducers.max=200;
set mapred.reduce.tasks= 200;---增大Reduce个数
set hive.groupby.mapaggr.checkinterval=100000 ;--这个是group的键对应的记录条数超过这个值则会进行分拆,值根据具体数据量设置
set hive.groupby.skewindata=true; --如果是group by过程出现倾斜 应该设置为true
set hive.skewjoin.key=100000; --这个是join的键对应的记录条数超过这个值则会进行分拆,值根据具体数据量设置
set hive.optimize.skewjoin=true;--如果是join 过程出现倾斜 应该设置为true
6:如果union all的部分个数大于2,或者每个union部分数据量大,应该拆成多个insert into 语句,实际测试过程中,执行时间能提升50%.
insert into table ..
select .. from (
   select .. from A
    union all 
   select .. from B
    union all 
   select .. from C
)
改为:
insert into table ..
select .. from A

insert into table ..
select .. from B

insert into table ..
select .. from C
7: 如何合并小文件,减少map数?
如果一个表中的map数特别多,可能是由于文件个数特别多,而且文件特别小照成的,可以进行如下操作,合并文件:
set mapred.max.split.size=100000000; // 100M
set mapred.min.split.size.per.node=100000000;
set mapred.min.split.size.per.rack=100000000;
set hive.input.format=org.apache.hadoop.hive.ql.io.CombineHiveInputFormat; // 合并小文件

但是,当map的业务逻辑很复杂时,需要释放增加map数。可以将table进行随机分布,使用新的table代替原表。
set mapred.reduce.tasks=10; 
create table temp as  
select * from a  
distribute by rand(123);
8: hive如何确定reduce数, reduce的个数基于以下参数设定:
hive.exec.reducers.bytes.per.reducer(每个reduce任务处理的数据量,默认为1000^3=1G)
hive.exec.reducers.max(每个任务最大的reduce数,默认为999)
计算reducer数的公式很简单N=min(参数2,总输入数据量/参数1)
即,如果reduce的输入(map的输出)总大小不超过1G,那么只会有一个reduce任务;所以调整以下参数:
set hive.exec.reducers.bytes.per.reducer=500000000; (500M)
set mapred.reduce.tasks = 15;

9:  Count(distinct)
当count distinct 的记录非常多的时候,设置以下两个参数:
hive.map.aggr = true
set hive.groupby.skewindata=true;

10: Group by
Group By的方法是在reduce做一些操作,这样会导致两个问题:
map端聚合,提前一部分计算:hive.map.aggr = true,
同时设置间隔:hive.groupby.mapaggr.checkinterval
均衡处理:hive.groupby.skewindata
这是针对数据倾斜的,设为ture的时候,任务的reduce会把原来一个job拆分成两个,第一个的job中reduce处理处理不同的随即分发过来的key的数据,生成中间结果,再由最后一个综合处理。

11: Order by, Sort by ,Dristribute by,Cluster By
order by VS Sort by: order by是在全局的排序,只用一个reduce去跑,所以在set hive.mapred.mode=strict 模式下,order by 必须limit,否则报错。Sort by只保证同一个reduce下排序正确。
Distribute by with sort by: Distribute by 是按指定的列把map 输出结果分配到reduce里。所以经常和sort by 来实现对某一字段的相同值分配到同一个reduce排序。
Cluster by 实现了Distribute by+ sort by 的功能.

12: 合并MapReduce操作
Multi-group by
Multi-group by是Hive的一个非常好的特性,它使得Hive中利用中间结果变得非常方便。例如:
FROM (SELECT a.status, b.school, b.gender

FROM status_updates a JOIN profiles b

ON (a.userid = b.userid and

a.ds='2009-03-20' )

) subq1

INSERT OVERWRITE TABLE gender_summary

PARTITION(ds='2009-03-20')

SELECT subq1.gender, COUNT(1) GROUP BY subq1.gender

INSERT OVERWRITE TABLE school_summary

PARTITION(ds='2009-03-20')

SELECT subq1.school, COUNT(1) GROUP BY subq1.school
上述查询语句使用了Multi-group by特性连续group by了2次数据,使用不同的group by key。这一特性可以减少一次MapReduce操作。

13: 参数调整
http://itindex.net/detail/53620-hive-%E4%BC%98%E5%8C%96

2015年7月14日星期二

大数据差异比较

问题:
有两份数据A和B,要求求解A和B的内容差异?
前提条件:
1. A和B使用文件保存,且数据量在G级别;
2. A和B的文件内容是结构化的;
3. A和B的文件内容是乱序的,即A的第一行可能在B中,但是在第n行,也可能不存在,B也一样。

解决方案:

1.linux diff命令

缺点:diff对文件进行按序比较,比如:
a.txt:
1 a
2 b
4 d
3 c
b.txt
2 b
4 d
3 c
1 a
5 e
执行命令:diff a.txt b.txt
1d0
< 1 a
4a4,5
> 1 a
> 5 e
分析结果可知,A和B同时包含的数据依然出现在结果中,需要对结果进行额外处理。当数据量较大时,diff运行超级慢。

2 分桶+grep
描述:A和B的数据是结构化的,因此可以取一列进行hash,将A和B进行分桶,比如分桶至128个文件中,然后分别比对A(1-128)和B(1-128)。
比对时,使用grep命令逐行扫两个文件即可。分桶可以使得grep的文件规模降低。
缺点:这是一个很笨的方法,速度及其慢。我采用了这个方法。呵呵!

3 sort+comm

sort命令是帮我们依据不同的数据类型进行排序,其语法及常用参数格式:
  sort [-bcfMnrtk][源文件][-o 输出文件]
补充说明:sort可针对文本文件的内容,以行为单位来排序。

参  数:
-b 忽略每行前面开始出的空格字符。
-c 检查文件是否已经按照顺序排序。
-C会检查文件是否已排好序,如果乱序,不输出内容,仅返回1。
-f 排序时,忽略大小写字母。
-M 将前面3个字母依照月份的缩写进行排序。
-n 依照数值的大小排序,防止10比2小的情况。
-o<输出文件> 将排序后的结果存入指定的文件。
-r 以相反的顺序来排序。
-t<分隔字符> 指定排序时所用的栏位分隔字符。
-k 选择以哪个区间进行排序。
-u 在输出行中去除重复行

例如,对a.txt进行排序:sort -n -k 1 -t ' ' a.txt

comm命令——对已经有序的文件进行比较
comm对文件进行处理时,要求文件已经有序,如果没有顺序,请使用sort进行排序后进行处理。语法:
comm [-123][--help][--version][第1个文件][第2个文件]

补充说明:
这项指令会一列列地比较两个已排序文件的差异,并将其结果显示出来,如果没有指定任何参数,则会把结果分成3行显示:
第1行仅是在第1个 文件中出现过的列;
第2行是仅在第2个文件中出现过的列;
第3行则是在第1与第2个文件里都出现过的列。

若给予的文件名称为"-",则comm指令会从标 准输入设备读取数据。
参数:
-1 不显示只在第1个文件里出现过的列。
-2 不显示只在第2个文件里出现过的列。
-3 不显示同时在第1和第2个文件里出现过的列。
例如:
comm a.txt b.txt 
               1 a
                2 b
                3 c
                4 d
        5 e
缺点:对于大文件,依然很慢。

4 egrep
egrep -f a.txt -v b.txt

5 使用hadoop和hive
最终采取的方案。
因为A和B都是结构化数据,因此可以使用hadoop和hive。创建两张表,将A和B分别载入,然后求内容差集。

使用hive的 关键字left semi join,解决的问题是:IN/EXISTS,即求并集。
例如:
select test_1.id, test_1.num from test_1 left semi join test_2 on (test_1.id = test_2.id);

使用hive的关键字left outer join, 解决A差B的问题:
例如:
select test_1.id, test_1.num from test_1 left outer join test_2 on (test_2.id = test_2.id) where test_2.num is null;

hive命令行配置条件

set mapred.job.priority=HIGH;
set mapred.job.groups=ods;
set mapred.job.queue.name=ods;
set mapred.map.tasks.speculative.execution = false;
set mapred.reduce.tasks.speculative.execution = false;
set mapred.job.map.capacity=1000;
set mapred.job.reduce.capacity=500;
set hive.metastore.client.socket.timeout=100000;
set hive.groupby.skewindata=true;
add jar /home/work/metastore.online/hive/lib/nova-pb-schema-1.0.jar;
add jar /home/work/metastore.online/hive/lib/hive-serde-2.3.33.jar;
add jar /home/work/metastore.online/hive/lib/hive-contrib-2.3.33.jar;

文件描述符

参考:http://www.bottomupcs.com/file_descriptors.html
文件描述符在形式上是一个非负整数。实际上,它是一个索引值,指向内核为每一个进程所维护的该进程打开文件的记录表。当程序打开一个现有文件或者创建一个新文件时,内核向进程返回一个文件描述符。在程序设计中,一些涉及底层的程序编写往往会围绕着文件描述符展开。但是文件描述符这一概念往往只适用于UNIX、Linux这样的操作系统。

每个进程在Linux内核中都有一个task_struct结构体来维护进程相关的 信息,称为进程描述符(Process Descriptor),而在操作系统理论中称为进程控制块 (PCB,Process Control Block)。task_struct中有一个指针(struct files_struct *files; )指向files_struct结构体,称为文件 描述符表,其中每个表项包含一个指向已打开的文件的指针。

用户程序不能直接访问内核中的文件描述符表,而只能使用文件描述符表的索引 (即0、1、2、3这些数字),这些索引就称为文件描述符(File Descriptor),用int 型变量保存。 当调用open 打开一个文件或创建一个新文件时,内核分配一个文件描述符并返回给用户程序,该文件描述符表项中的指针指向新打开的文件。当读写文件时,用户程序把文件描述符传给read 或write ,内核根据文件描述符找到相应的表项,再通过表项中的指针找到相应的文件。
Java对应的文件描述符类为:FileDescriptor
使用Demo:
// 新建FileInputStream对象
File file = new File(FileName);
FileInputStream in1 = new FileInputStream(file);
// 获取文件“file.txt”对应的“文件描述符”
FileDescriptor fdin = in2.getFD();
// 根据“文件描述符”创建“FileInputStream”对象
 FileInputStream in3 = new FileInputStream(fdin);

2015年7月12日星期日

shell学习笔记-输入输出,重复执行命令,字段分隔符,比较与测试

一 输入输出
子shell通过()操作符进行定义,子shell的改变不会影响主shell
ps:通过引用子shell的方式保留空格和换行符时,需要使用双引号 "(..)"
#管道
cmd1 | cmd2   #e.g.   ls | cat -n >out.txt
#子shell
cmd_output=$(ls | cat -n)
#反引用
cmd_output=`cmd`
二 运行命令直至执行成功

repeat(){while :; do &@ && return;  sleep 20; done } 
#&@表示传入的命令  例如
repeat wget -c http://....
true作为二进制文件实现,因此使用:命令,:命令默认返回为0的退出码

三 字段分隔符和迭代器

内部字段分隔符为IFS
data="one,two,three,four"
old_ifs=$IFS
IFS=,
for item in $data;
do
        echo $item
done
IFS=$old_ifs
四 比较与测试
-a逻辑与    condition1 -a condition2 
if [[  condition1 ]] && [[ condition2 ]]
-o逻辑或    condition1 -a condition2  
if [[  condition1 ]] || [[ condition2 ]]

shell学习笔记-文件描述符,数组,别名,调试

由于工作中需要写shell,因此进行系统性学习,并记录学习笔记。使用的学习资料为《Linux Shell脚本攻略》第二版。

一 文件描述符

 自定义文件描述符:

exec 3<input.txt #使用文件进行文件描述符输入
echo "string" >&4 #写入文件描述符4
cat<&3  #读取文件描述符3
二 数组和关联数组(Map)

数组: 
#定义
arr=(1 2 3 4)

echo "all:"
#打印所有元素
echo ${arr[*]}

echo ${arr[@]}
#打印长度
echo "length:" ${#arr[*]}
关联数组:
#deifination
declare -A ass_arr
#assignment
ass_arr=(['apple']=10 ['pear']=20)
echo ${ass_arr['apple']}
#get all keys
echo ${!ass_arr[*]}
三 别名
alias cmd='new cmd'
使用\cmd 进行转义,可以忽略别名

四 调试
#全局调试
sh  -x script 或者 #!/bin/bash -vx
#局部调试
set -x
cmd
set +x

2015年7月10日星期五

Disruptor详解

对Disruptor的最初印象就是ringbuffer。但是尽管ringbuffer是整个模式(Disruptor)的核心,但是Disruptor对ringbuffer的访问控制策略才是真正的关键点所在。
  • ringbuffer到底是什么?
正如名字所说的一样,它是一个环(首尾相接的环),你可以把它用做在不同上下文(线程)间传递数据的buffer。




基本来说,ringbuffer拥有一个序号,这个序号指向数组中下一个可用的元素。(校对注:如下图右边的图片表示序号,这个序号指向数组的索引4的位置。)



随着你不停地填充这个buffer(可能也会有相应的读取),这个序号会一直增长,直到绕过这个环。



要找到数组中当前序号指向的元素,可以通过mod操作:

sequence mod array length = array index

以上面的ringbuffer为例(java的mod语法):12 % 10 = 2。

事实上,上图中的ringbuffer只有10个槽完全是个意外。如果槽的个数是2的N次方更有利于基于二进制的计算机进行计算。

(校对注:2的N次方换成二进制就是1000,100,10,1这样的数字, sequence & (array length-1) = array index,比如一共有8槽,3&(8-1)=3,HashMap就是用这个方式来定位数组元素的,这种方式比取模的速度更快。)
那又怎么样?

如果你看了维基百科里面的关于环形buffer的词条,你就会发现,ringbuffer的实现方式,与其最大的区别在于:没有尾指针。ringbuffer只维护了一个指向下一个可用位置的序号。这种实现是经过深思熟虑的—ringbuffer选择用环形buffer的最初原因就是想要提供可靠的消息传递。我们需要将已经被服务发送过的消息保存起来,这样当另外一个服务通过nak (校对注:拒绝应答信号)告诉ringbuffer没有成功收到消息时,ringbuffer能够重新发送给他们。

听起来,环形buffer非常适合这个场景。它维护了一个指向尾部的序号,当收到nak(校对注:拒绝应答信号)请求,可以重发从那一点到当前序号之间的所有消息:



ring buffer和大家常用的队列之间的区别是,ringbuffer不删除buffer中的数据,也就是说这些数据一直存放在buffer中,直到新的数据覆盖他们。这就是和维基百科版本相比,ringbuffer不需要尾指针的原因。ringbuffer本身并不控制是否需要重叠(决定是否重叠是生产者-消费者行为模式的一部分
  • 它为什么如此优秀?
之所以ringbuffer采用这种数据结构,是因为它在可靠消息传递方面有很好的性能。这就够了,不过它还有一些其他的优点。

首先,因为它是数组,所以要比链表快,而且有一个容易预测的访问模式。(译者注:数组内元素的内存地址的连续性存储的)。这是对CPU缓存友好的—也就是说,在硬件级别,数组中的元素是会被预加载的,因此在ringbuffer当中,cpu无需时不时去主存加载数组中的下一个元素。(校对注:因为只要一个元素被加载到缓存行,其他相邻的几个元素也会被加载进同一个缓存行)

其次,你可以为数组预先分配内存,使得数组对象一直存在(除非程序终止)。这就意味着不需要花大量的时间用于垃圾回收。此外,不像链表那样,需要为每一个添加到其上面的对象创造节点对象—对应的,当删除节点时,需要执行相应的内存清理操作。

  • ringbuffer的为什么这么快?--cache line padding
设想你的long类型的数据不是数组的一部分。设想它只是一个单独的变量。让我们称它为head,这么称呼它其实没有什么原因。然后再设想在你的类中有另一个变量紧挨着它。让我们直接称它为tail。现在,当你加载head到缓存的时候,你也免费加载了tail。

听想来不错。直到你意识到tail正在被你的生产者写入,而head正在被你的消费者写入。这两个变量实际上并不是密切相关的,而事实上却要被两个不同内核中运行的线程所使用。
设想你的消费者更新了head的值。缓存中的值和内存中的值都被更新了,而其他所有存储head的缓存行都会都会失效,因为其它缓存中head不是最新值了。

现在如果一些正在其他内核中运行的进程只是想读tail的值,整个缓存行需要从主内存重新读取。那么一个和你的消费者无关的线程读一个和head无关的值,它被缓存未命中给拖慢了。

当然如果两个独立的线程同时写两个不同的值会更糟。因为每次线程对缓存行进行写操作时,每个内核都要把另一个内核上的缓存块无效掉并重新读取里面的数据。你基本上是遇到两个线程之间的写冲突了,尽管它们写入的是不同的变量。

这叫作“伪共享”(译注:可以理解为错误的共享),因为每次你访问head你也会得到tail,而且每次你访问tail,你也会得到head。这一切都在后台发生,并且没有任何编译警告会告诉你,你正在写一个并发访问效率很低的代码。

解决方案-神奇的缓存行填充

你会看到Disruptor消除这个问题,至少对于缓存行大小是64字节或更少的处理器架构来说是这样的(译注:有可能处理器的缓存行是128字节,那么使用64字节填充还是会存在伪共享问题),通过增加补全来确保ring buffer的序列号不会和其他东西同时存在于一个缓存行中。
1public long p1, p2, p3, p4, p5, p6, p7; // cache line padding
2    private volatile long cursor = INITIAL_CURSOR_VALUE;
3    public long p8, p9, p10, p11, p12, p13, p14; // cache line padding


因此没有伪共享,就没有和其它任何变量的意外冲突,没有不必要的缓存未命中。
在你的Entry类中也值得这样做,如果你有不同的消费者往不同的字段写入,你需要确保各个字段间不会出现伪共享。

整理自:http://ifeve.com/disruptor-writing-ringbuffer/