레이블이 Hadoop인 게시물을 표시합니다. 모든 게시물 표시
레이블이 Hadoop인 게시물을 표시합니다. 모든 게시물 표시

2012년 2월 22일 수요일

Pig UDFs


lzo 파일을 pig에서 사용해야 해서 UDF를 만들어 봤다. 
 -  hadoop에 lzo 설정하기 
 


이전까지 많은 작업들이 약간은(?) 복잡한 알고리즘이 필요한 작업들이라서 java 로 직접 데이터를 만졌었는데, 반복적인 작업들이 필요하게 되면서 pig로 작업을 하는게 편할것 같았다. 





필요한 udf 또한 복잡하지 않고 너무도 간단한 수준이라서 udf에 대해서 자세한 내용은 모르지만 필요한 함수를 만들어 보았다. (elephant-bird 라고 트위터에서 오픈해 놓은 코드가 있는데, 사내에서 사용하기가 어려운 상황..)
 


jython을 이용한 udf는 만들기가 간단해서 활용도가 높을것 같다. 



# Java UDFs - LzoPigStorage  


package xxxxx;


 


import org.apache.hadoop.io.LongWritable;


import org.apache.hadoop.io.Text;


import org.apache.hadoop.mapreduce.InputFormat;


import org.apache.hadoop.mapreduce.OutputFormat;


import org.apache.hadoop.mapreduce.lib.output.TextOutputFormat;


import org.apache.pig.builtin.PigStorage;


import com.hadoop.mapreduce.LzoTextInputFormat;


 


public class LzoPigStorage extends PigStorage {


private String delimiter = null;


 


public LzoPigStorage() {


super();


}


 


public LzoPigStorage(String delimiter) {


super(delimiter);


this.delimiter = delimiter;


}


 


@Override


public InputFormat<LongWritable, Text> getInputFormat() {


return new LzoTextInputFormat();


}


@Override


        public OutputFormat getOutputFormat() {


            return new TextOutputFormat();


        }


}


 


// 사용


register 파일명.jar;


A = load 'data_path' using xxx.LzoPigStorage('\t') AS (.....);





 





# jython 









// string_pig_udf.py





@outputSchema("rquery:chararray")



def rmQuerySpace(instr):



    return instr.replace(' ','')








// 사용


register 'string_pig_udf.py' using jython as myfuncs;


...


C = FOREACH B GENERATE myfuncs.reQuerySpace(query);


 












# 간단히 wiki로도 정리  





2011년 7월 19일 화요일

Hadoop Lzo 압축 설정 (2)


이전에 Hadoop Lzo 압축 설정에 관한 블로깅을 했는데...

그 이후에 발생했던 문제와,
실질적으로 사용하는 방법에 대해서 정리..


# 두 가지 버전의 hadoop lzo
1. https://github.com/omalley/hadoop-gpl-compression
  - 이전에 설치했던 버전

2. https://github.com/kevinweil/hadoop-lzo
  - 이 버전으로 다시 설치


# 사용법

1. 파일시스템의 파일을 압축해서 hdfs에 올리는 방법.
> lzop 파일이름
> hadoop fs -copyFromLocal 파일이름.lzo hdfs위치



2. hdfs의 lzo 파일에 index 만들기
 - lzo 파일이 있는곳에 파일이름.index라는 파일이 생긴다.
 - index를 안해주면, lzo 파일을 split하지 않고 하나의 map으로 처리

1) index it in-process via:
hadoop jar /path/to/your/hadoop-lzo.jar com.hadoop.compression.lzo.LzoIndexer big_file.lzo

2) index it in a map-reduce job via:
hadoop jar /path/to/your/hadoop-lzo.jar com.hadoop.compression.lzo.DistributedLzoIndexer big_file.lzo




3. splitted lzo 파일 사용
  - 파일이름.index 를 보고 알아서 나누어서 작업한다.
1) java
job 설정에 다음을 추가..

job.setInputFormatClass(LzoTextInputFormat.class);
 
2) streaming (테스트 안해봄)
실행할때 다음을 추가

"-inputformat com.hadoop.mapred.DeprecatedLzoTextInputFormat




4. 최종 결과 파일을 fs로
1) lzo_deflate
  - 압축코덱을 LzoCodec으로 하면 .lzo_deflate의 형태로 압축된다.
  - lzo_deflate는 파일시스템으로 getmerge 한 후 압축을 어떻게 푸는지 알수가 없어서, 코덱 설정을 바꿈.


<property>
      <name>mapred.output.compression.codec</name>
      <value>com.hadoop.compression.lzo.LzoCodec</value>
</property>


2) LzoCodec -> LzopCodec
   - The LzoCodec for the pure LZO
format, which uses the .lzo_deflate filename extension (by analogy with
DEFLATE, which is gzip without the headers).
  - The LzopCodec is compatible with the lzop tool, which is essentially the LZO format with extra headers.
 - map 과 reduce 사이에는 header가 없는 LzoCodec으로..
 - 최종 output은 LzopCodec 으로 변경


<property>
      <name>mapred.output.compression.codec</name>
     <value>com.hadoop.compression.lzo.LzopCodec</value>
</property>



3) getmerge 후에 fs에서 압출 풀때
  - 1개의 reducer의 경우 'lzop -d' 로 일반적인 경우처럼
  - reducer의 개수 1개 이상일 때는 다음과 같이
lzop -d 파일명.lzo -o 압축풀파일명


2011년 7월 6일 수요일

Hadoop LZO 압축 설정


# 첨가. 

hadoop-gpl-compression 를 

https://github.com/kevinweil/hadoop-lzo 로 변경해서 설치할것. 

참고 : http://upepo.tistory.com/174

-------------------------



1. lzo 설치


# 다운로드
http://www.oberhumer.com/opensource/lzo/

# 설치
./configure --enable-shared
make
make instll

# 기타
LD_LIBRARY_PATH 설정
or
/sbin/ldconfig



2. native connector library 설치

# 다운로드
http://code.google.com/a/apache-extras.org/p/hadoop-gpl-compression/

# hadoop library 복사
hadoop.0.20.0-core.jar 를 hadoop-gpl-compression/lib/ 으로 복사

# build
cnt compile-native
ant jar



3. 설정..

# 64bit 인 경우
다음 파일들을..
hadoop-gpl-compression/build/native/Linux-amd64-64/libgplcompression.la
hadoop-gpl-compression/build/native/Linux-amd64-64/lib/*

여기로 복사
hadoop/lib/native/Linux-amd64-64/

다음 파일을
hadoop-gpl-compression/build/hadoop-gpl-compression-0.1.0-dev.jar

여기로 복사
hadoop/lib


# 32bit 인 경우
Linux-amd64-64 --> Linux-i386-32 로 해서 위와 같게..



4. .bash_profile 에 추가

JAVA_LIBRARY_PATH=$JAVA_LIBRRAY_PATH:$HADOOP_HOME/lib/native/Linux-amd64-64/
JAVA_LIBRARY_PATH=$JAVA_LIBRRAY_PATH:$HADOOP_HOME/lib/

export JAVA_LIBRARY_PATH




5. lzop 설치

http://www.lzop.org/
필요하면 설치..




6. hadoop conf

# core-site.xml
- lzo compression도 기본 코덱으로 설정
        <property>
              <name>io.compression.codecs</name>
        <value>org.apache.hadoop.io.compress.GzipCodec,

org.apache.hadoop.io.compress.DefaultCodec,


org.apache.hadoop.io.compress.BZip2Codec,


com.hadoop.compression.lzo.LzoCodec</value>

        </property>

        <property>
                <name>io.compression.codec.lzo.class</name>
                <value>com.hadoop.compression.lzo.LzoCodec</value>
        </property>




# mapred-site.xml
- map이 끝나고 압축해서 reduce로 전달
        <property>
                <name>mapred.compress.map.output</name>
                <value>true</value>
        </property>
        <property>
                <name>mapred.map.output.compression.codec</name>
                <value>com.hadoop.compression.lzo.LzoCodec</value>
        </property>

- reduce 끝나고 최종 결과 압축
        <property>
                <name>mapred.output.compress</name>
                <value>true</value>
        </property>
        <property>
                <name>mapred.output.compression.codec</name>
                <value>com.hadoop.compression.lzo.LzoCodec</value>
        </property>



...

설치 중에 문제 생기는 부분이 있으면 답글 남겨 주세요~ 


2011년 5월 16일 월요일

Python 하둡 스트리밍 (Hadoop Streaming) #2

이전글: Python 하둡 스트리밍 (Hadoop Streamming) #1 

참조
[1] : https://github.com/jhofman/icwsm2010_tutorial/blob/master/hstream.py
[2] : http://www.michael-noll.com/tutorials/writing-an-hadoop-mapreduce-program-in-python
[3] : http://jakehofman.com/icwsm2010 



파이썬으로 hadoop streaming을 편하게 할수 있는 파이썬 클래스 소개


# 실행


./bin/hadoop jar contrib/streaming/hadoop-0.20.2-streaming.jar \


    -file .../wordcount.py \
    -file .../hstream.py \
    -mapper '.../wordcount.py -m'  \


    -reducer '.../wordcount.py -r' \


    -input input_data  \


    -output output_data 




# wordcount.py


#!/usr/bin/env python



from hstream import HStream


import sys


import re


from collections import defaultdict



class WordCount(HStream):


        def mapper(self, record):


                for word in " ".join(record).split():


                        self.write_output((word,1))



        def reducer(self, key, records):


                total = 0


                for record in records:


                        word, count = record


                        total += int(count)


                self.write_output((word,total))



if __name__== '__main__':


        WordCount()


2011년 4월 18일 월요일

Python 하둡 스트리밍 (Hadoop Streaming) #1


# 참조
http://www.michael-noll.com/tutorials/writing-an-hadoop-mapreduce-program-in-python/


# 실행


./bin/hadoop jar contrib/streaming/hadoop-0.20.2-streaming.jar \


    -file .../mapper.py -mapper .../mapper.py  \


    -file .../reducer.py -reducer .../reducer.py  \


    -input input_data  \


    -output output_data






# 기본 예제 : word count

1. mapper.py


#!/usr/bin/env python


import sys


 


for line in sys.stdin:


        line = line.strip()


        words = line.split()


        for word in words:


                print '%s\t%s' % (word,1)




2. reducer.py


#!/usr/bin/env python


import sys


from operator import itemgetter
 



word2count = {} 


for line in sys.stdin:


        line = line.strip()


        word, count = line.split('\t', 1)


        try:


                count = int(count)


                word2count[word] = word2count.get(word, 0) + count


        except:


                pass


# sort


sorted_word2count = sorted(word2count.items(), key=itemgetter(0))


 


for word, count in sorted_word2count:


        print '%s\t%s' % (word, count)







기본 예제를 조금만 수정하면,
간단한 처리는 대부분 가능하다. 

#
참조 링크에 yield 와 groupby를 이용하는 코드가 나오는데 python 스럽기도 하고,
실제도 돌려보니 속도도 조금 빠르다.