Hadoop Streaming
Hadoop Streaming
Hadoop 배포에 포함된 유틸리티로, 임의의 실행 파일이나 스크립트를 매퍼·리듀서로 사용해 Map/Reduce 작업을 만들고 실행할 수 있게 하는 Hadoop streaming을 설명하는 문서예요. 동작 방식, 명령 옵션, 예시, FAQ를 다룹니다.
출처: 문서
본문
Hadoop Streaming
Hadoop streaming은 Hadoop 배포에 포함된 유틸리티입니다. 이 유틸리티는 임의의 실행 파일이나 스크립트를 매퍼·리듀서로 사용해 Map/Reduce 작업을 만들고 실행할 수 있게 해 줍니다. 예:
mapred streaming \
-input myInputDirs \
-output myOutputDir \
-mapper /bin/cat \
-reducer /usr/bin/wc
streaming 동작 방식 (How Streaming Works)
위 예시에서 매퍼와 리듀서는 모두 stdin에서 입력을 읽고(줄 단위로) stdout으로 출력을 내보내는 실행 파일입니다. 유틸리티는 Map/Reduce 작업을 만들고, 적절한 클러스터에 작업을 제출하며, 작업이 완료될 때까지 진행을 모니터링합니다.
매퍼에 실행 파일이 지정되면, 각 매퍼 작업은 매퍼가 초기화될 때 실행 파일을 별도 프로세스로 실행합니다. 매퍼 작업이 실행되면서 입력을 줄로 변환해 프로세스의 stdin에 공급합니다. 그 사이 매퍼는 프로세스의 stdout에서 줄 지향 출력을 수집하고 각 줄을 키/값 쌍으로 변환해 매퍼의 출력으로 수집합니다. 기본적으로 줄의 첫 번째 탭 문자까지의 접두어 가 key이고, 줄의 나머지(탭 문자 제외)가 value입니다. 줄에 탭 문자가 없으면 전체 줄이 키로 간주되고 값은 null입니다. 하지만 아래에서 설명하는 것처럼 -inputformat 명령 옵션을 설정해 사용자 정의할 수 있습니다.
리듀서에 실행 파일이 지정되면, 각 리듀서 작업은 리듀서가 초기화될 때 실행 파일을 별도 프로세스로 실행합니다. 리듀서 작업이 실행되면서 입력 키/값 쌍을 줄로 변환해 프로세스의 stdin에 공급합니다. 그 사이 리듀서는 프로세스의 stdout에서 줄 지향 출력을 수집하고 각 줄을 키/값 쌍으로 변환해 리듀서의 출력으로 수집합니다. 기본적으로 줄의 첫 번째 탭 문자까지의 접두어가 키이고, 줄의 나머지(탭 문자 제외)가 값입니다. 아래에서 설명하는 -outputformat 명령 옵션으로 사용자 정의할 수 있어요.
이것이 Map/Reduce 프레임워크와 streaming 매퍼/리듀서 사이의 통신 프로토콜의 기초입니다.
사용자는 stream.non.zero.exit.is.failure를 true 또는 false로 지정해 0이 아닌 상태로 종료하는 streaming 작업을 각각 Failure 또는 Success로 만들 수 있습니다. 기본적으로 0이 아닌 상태로 종료하는 streaming 작업은 실패한 작업으로 간주됩니다.
Streaming 명령 옵션 (Streaming Command Options)
Streaming은 streaming 명령 옵션과 generic 명령 옵션을 모두 지원합니다. 일반 명령줄 문법은 아래와 같습니다.
참고: generic 옵션을 streaming 옵션 앞에 배치해야 합니다. 그렇지 않으면 명령이 실패할 수 있어요. 예시는 Making Archives Available to Tasks를 참고하세요.
mapred streaming [genericOptions] [streamingOptions]
Hadoop streaming 명령 옵션은 다음과 같습니다.
| 파라미터 | 선택/필수 | 설명 |
|---|---|---|
| -input directoryname or filename | 필수 | 매퍼용 입력 위치 |
| -output directoryname | 필수 | 리듀서용 출력 위치 |
| -mapper executable or JavaClassName | 선택 | 매퍼 실행 파일. 지정하지 않으면 IdentityMapper가 기본. |
| -reducer executable or JavaClassName | 선택 | 리듀서 실행 파일. 지정하지 않으면 IdentityReducer가 기본. |
| -file filename | 선택 | 매퍼·리듀서·컴바이너 실행 파일을 컴퓨트 노드에 로컬로 사용 가능하게 함. |
| -inputformat JavaClassName | 선택 | 제공하는 클래스는 Text 클래스의 키/값 쌍을 반환해야 함. 지정하지 않으면 TextInputFormat이 기본. |
| -outputformat JavaClassName | 선택 | 제공하는 클래스는 Text 클래스의 키/값 쌍을 받아야 함. 지정하지 않으면 TextOutputformat이 기본. |
| -partitioner JavaClassName | 선택 | 키를 어떤 리듀스로 보낼지 결정하는 클래스. |
| -combiner streamingCommand or JavaClassName | 선택 | map 출력용 컴바이너 실행 파일. |
| -cmdenv name=value | 선택 | environment 변수를 streaming 명령에 전달. |
| -inputreader | 선택 | 하위 호환용: 기록 리더 클래스(입력 포맷 클래스 대신) 지정. |
| -verbose | 선택 | 자세한 출력. |
| -lazyOutput | 선택 | 지연 생성 출력. 예: 출력 포맷이 FileOutputFormat 기반이면 Context.write 첫 호출에서만 출력 파일 생성. |
| -numReduceTasks | 선택 | 리듀서 수 지정. |
| -mapdebug | 선택 | map 작업 실패 시 호출할 스크립트. |
| -reducedebug | 선택 | reduce 작업 실패 시 호출할 스크립트. |
매퍼/리듀서로 Java 클래스 지정
매퍼·리듀서로 Java 클래스를 제공할 수 있습니다.
mapred streaming \
-input myInputDirs \
-output myOutputDir \
-inputformat org.apache.hadoop.mapred.KeyValueTextInputFormat \
-mapper org.apache.hadoop.mapred.lib.IdentityMapper \
-reducer /usr/bin/wc
stream.non.zero.exit.is.failure를 true 또는 false로 지정해 0이 아닌 상태로 종료하는 streaming 작업을 각각 Failure 또는 Success로 만들 수 있습니다. 기본적으로 0이 아닌 상태로 종료하는 streaming 작업은 실패로 간주됩니다.
작업 제출 시 파일 패키징 (Packaging Files With Job Submissions)
임의의 실행 파일을 매퍼·리듀서로 지정할 수 있습니다. 실행 파일이 클러스터 머신에 미리 존재할 필요는 없지만, 없다면 "-file" 옵션으로 프레임워크에 실행 파일을 작업 제출의 일부로 패키징하라고 알려야 합니다. 예:
mapred streaming \
-input myInputDirs \
-output myOutputDir \
-mapper myPythonScript.py \
-reducer /usr/bin/wc \
-file myPythonScript.py
위 예시는 사용자 정의 Python 실행 파일을 매퍼로 지정합니다. "-file myPythonScript.py" 옵션은 python 실행 파일이 작업 제출의 일부로 클러스터 머신에 배송되게 합니다.
실행 파일 외에도 매퍼·리듀서가 사용할 수 있는 다른 보조 파일(사전, 구성 파일 등)도 패키징할 수 있습니다. 예:
mapred streaming \
-input myInputDirs \
-output myOutputDir \
-mapper myPythonScript.py \
-reducer /usr/bin/wc \
-file myPythonScript.py \
-file myDictionary.txt
작업에 다른 플러그인 지정
일반 Map/Reduce 작업과 마찬가지로 streaming 작업에 다른 플러그인을 지정할 수 있습니다.
-inputformat JavaClassName
-outputformat JavaClassName
-partitioner JavaClassName
-combiner streamingCommand or JavaClassName
입력 포맷으로 제공하는 클래스는 Text 클래스의 키/값 쌍을 반환해야 합니다. 입력 포맷 클래스를 지정하지 않으면 TextInputFormat이 기본으로 사용됩니다. TextInputFormat은 실제로 입력 데이터의 일부가 아닌 LongWritable 클래스의 키를 반환하므로 키는 버려지고 값만 streaming 매퍼로 파이프됩니다.
출력 포맷으로 제공하는 클래스는 Text 클래스의 키/값 쌍을 받을 것으로 기대됩니다. 출력 포맷 클래스를 지정하지 않으면 TextOutputFormat이 기본으로 사용됩니다.
환경 변수 설정
streaming 명령에서 환경 변수를 설정하려면:
-cmdenv EXAMPLE_DIR=/home/example/dictionaries/
Generic 명령 옵션 (Generic Command Options)
Streaming은 streaming 명령 옵션과 generic 명령 옵션을 모두 지원합니다. 일반 명령줄 문법은 아래와 같습니다.
참고: generic 옵션을 명령 옵션 앞에 배치해야 합니다. 그렇지 않으면 명령이 실패할 수 있어요. 예시는 Making Archives Available to Tasks를 참고하세요.
hadoop command [genericOptions] [streamingOptions]
streaming과 함께 쓸 수 있는 Hadoop generic 명령 옵션은 다음과 같습니다.
| 파라미터 | 선택/필수 | 설명 |
|---|---|---|
| -conf configuration_file | 선택 | 애플리케이션 구성 파일 지정 |
| -D property=value | 선택 | 주어진 속성에 value 사용 |
| -fs host:port or local | 선택 | 네임노드 지정 |
| -files | 선택 | Map/Reduce 클러스터로 복사할 파일들을 쉼표 구분으로 지정 |
| -libjars | 선택 | 클래스패스에 포함할 jar 파일들을 쉼표 구분으로 지정 |
| -archives | 선택 | 컴퓨트 머신에서 압축 해제할 아카이브들을 쉼표 구분으로 지정 |
-D 옵션으로 구성 변수 지정
"-D
디렉터리 지정
로컬 임시 디렉터리를 바꾸려면:
-D dfs.data.dir=/tmp
추가 로컬 임시 디렉터리를 지정하려면:
-D mapred.local.dir=/tmp/local
-D mapred.system.dir=/tmp/system
-D mapred.temp.dir=/tmp/temp
참고: 작업 구성 파라미터의 더 자세한 내용은 mapred-default.xml 참고.
Map-Only 작업 지정
입력 데이터를 map 함수만으로 처리하고 싶을 때가 많습니다. 이렇게 하려면 mapreduce.job.reduces를 0으로 설정하면 됩니다. Map/Reduce 프레임워크는 리듀서 작업을 만들지 않습니다. 대신 매퍼 작업의 출력이 작업의 최종 출력이 됩니다.
-D mapreduce.job.reduces=0
하위 호환을 위해 Hadoop Streaming은 "-reducer NONE" 옵션도 지원하며, 이것은 "-D mapreduce.job.reduces=0"과 동일합니다.
리듀서 수 지정
리듀서 수를 지정하려면, 예를 들어 2개는:
mapred streaming \
-D mapreduce.job.reduces=2 \
-input myInputDirs \
-output myOutputDir \
-mapper /bin/cat \
-reducer /usr/bin/wc
줄을 키/값 쌍으로 나누는 방식 사용자 정의
앞서 언급했듯 Map/Reduce 프레임워크는 매퍼의 stdout에서 줄을 읽으면 그 줄을 키/값 쌍으로 나눕니다. 기본적으로 첫 번째 탭 문자까지의 접두어가 키이고, 나머지(탭 문자 제외)가 값입니다.
하지만 이 기본값을 사용자 정의할 수 있습니다. 탭 문자(기본)가 아닌 다른 필드 구분자를 지정할 수 있고, 줄의 첫 번째 문자(기본)가 아닌 n번째(n >= 1) 문자를 키와 값의 구분자로 지정할 수 있습니다. 예:
mapred streaming \
-D stream.map.output.field.separator=. \
-D stream.num.map.output.key.fields=4 \
-input myInputDirs \
-output myOutputDir \
-mapper /bin/cat \
-reducer /bin/cat
위 예시에서 "-D stream.map.output.field.separator=."은 map 출력의 필드 구분자로 "."을 지정하고, 줄의 네 번째 "."까지의 접두어가 키이고 나머지(네 번째 "." 제외)가 값입니다. 줄에 "."이 네 개 미만이면 전체 줄이 키이고 값은 빈 Text 객체(new Text("")로 만든 것)입니다.
마찬가지로 "-D stream.reduce.output.field.separator=SEP"와 "-D stream.num.reduce.output.fields=NUM"으로 reduce 출력의 줄에서 n번째 필드 구분자를 키와 값의 구분자로 지정할 수 있습니다.
마찬가지로 Map/Reduce 입력의 입력 구분자로 "stream.map.input.field.separator"와 "stream.reduce.input.field.separator"를 지정할 수 있습니다. 기본 구분자는 탭 문자입니다.
큰 파일과 아카이브 작업 (Working with Large Files and Archives)
-files와 -archives 옵션은 파일·아카이브를 작업에 사용 가능하게 합니다. 인자는 이미 HDFS에 업로드한 파일·아카이브의 URI입니다. 이 파일·아카이브는 작업 간에 캐시됩니다. host와 fs_port 값을 fs.default.name 구성 변수에서 가져올 수 있습니다.
참고: -files와 -archives 옵션은 generic 옵션입니다. 반드시 generic 옵션을 명령 옵션 앞에 배치해야 합니다. 그렇지 않으면 명령이 실패할 수 있어요.
파일을 작업에 사용 가능하게 하기
-files 옵션은 파일의 로컬 복사본을 가리키는 심볼릭 링크를 작업의 현재 작업 디렉터리에 만듭니다.
이 예시에서 Hadoop은 작업의 현재 작업 디렉터리에 testfile.txt라는 심볼릭 링크를 자동으로 만듭니다. 이 심볼릭 링크는 testfile.txt의 로컬 복사본을 가리킵니다.
-files hdfs://host:fs_port/user/testfile.txt
사용자는 #으로 -files에 다른 심볼릭 링크 이름을 지정할 수 있습니다.
-files hdfs://host:fs_port/user/testfile.txt#testfile
여러 항목은 이렇게 지정할 수 있습니다.
-files hdfs://host:fs_port/user/testfile1.txt,hdfs://host:fs_port/user/testfile2.txt
아카이브를 작업에 사용 가능하게 하기
-archives 옵션은 jar를 작업의 현재 작업 디렉터리에 로컬로 복사하고 파일을 자동으로 압축 해제하게 합니다.
이 예시에서 Hadoop은 작업의 현재 작업 디렉터리에 testfile.jar라는 심볼릭 링크를 자동으로 만듭니다. 이 심볼릭 링크는 업로드된 jar 파일의 압축 해제된 내용을 저장하는 디렉터리를 가리킵니다.
-archives hdfs://host:fs_port/user/testfile.jar
사용자는 #으로 -archives에 다른 심볼릭 링크 이름을 지정할 수 있습니다.
-archives hdfs://host:fs_port/user/testfile.tgz#tgzdir
이 예시에서 input.txt 파일은 두 파일(cachedir.jar/cache.txt와 cachedir.jar/cache2.txt)의 이름을 지정하는 두 줄을 가집니다. "cachedir.jar"는 "cache.txt"와 "cache2.txt" 파일을 가진 아카이브된 디렉터리의 심볼릭 링크입니다.
mapred streaming \
-archives 'hdfs://hadoop-nn1.example.com/user/me/samples/cachefile/cachedir.jar' \
-D mapreduce.job.maps=1 \
-D mapreduce.job.reduces=1 \
-D mapreduce.job.name="Experiment" \
-input "/user/me/samples/cachefile/input.txt" \
-output "/user/me/samples/cachefile/out" \
-mapper "xargs cat" \
-reducer "cat"
$ ls test_jar/
cache.txt cache2.txt
$ jar cvf cachedir.jar -C test_jar/ .
added manifest
adding: cache.txt(in = 30) (out= 29)(deflated 3%)
adding: cache2.txt(in = 37) (out= 35)(deflated 5%)
$ hdfs dfs -put cachedir.jar samples/cachefile
$ hdfs dfs -cat /user/me/samples/cachefile/input.txt
cachedir.jar/cache.txt
cachedir.jar/cache2.txt
$ cat test_jar/cache.txt
This is just the cache string
$ cat test_jar/cache2.txt
This is just the second cache string
$ hdfs dfs -ls /user/me/samples/cachefile/out
Found 2 items
-rw-r--r-* 1 me supergroup 0 2013-11-14 17:00 /user/me/samples/cachefile/out/_SUCCESS
-rw-r--r-* 1 me supergroup 69 2013-11-14 17:00 /user/me/samples/cachefile/out/part-00000
$ hdfs dfs -cat /user/me/samples/cachefile/out/part-00000
This is just the cache string
This is just the second cache string
더 많은 사용 예시 (More Usage Examples)
Hadoop Partitioner 클래스
Hadoop에는 많은 애플리케이션에 유용한 라이브러리 클래스 KeyFieldBasedPartitioner가 있습니다. 이 클래스는 Map/Reduce 프레임워크가 전체 키가 아니라 특정 키 필드에 기반해 map 출력을 파티셔닝하게 합니다. 예:
mapred streaming \
-D stream.map.output.field.separator=. \
-D stream.num.map.output.key.fields=4 \
-D map.output.key.field.separator=. \
-D mapreduce.partition.keypartitioner.options=-k1,2 \
-D mapreduce.job.reduces=12 \
-input myInputDirs \
-output myOutputDir \
-mapper /bin/cat \
-reducer /bin/cat \
-partitioner org.apache.hadoop.mapred.lib.KeyFieldBasedPartitioner
여기서 -D stream.map.output.field.separator=. 와 -D stream.num.map.output.key.fields=4 는 이전 예시에서 설명한 대로입니다. 두 변수는 streaming이 매퍼의 키/값 쌍을 식별하는 데 사용합니다.
위 Map/Reduce 작업의 map 출력 키는 보통 "."으로 구분된 네 필드를 가집니다. 하지만 Map/Reduce 프레임워크는 -D mapred.text.key.partitioner.options=-k1,2 옵션으로 키의 처음 두 필드로 map 출력을 파티셔닝합니다. 여기서 -D map.output.key.field.separator=. 는 파티션 구분자를 지정합니다. 이는 키의 처음 두 필드가 같은 모든 키/값 쌍이 같은 리듀서로 파티셔닝되도록 보장합니다.
이것은 사실상 처음 두 필드를 기본 키로, 다음 두 필드를 보조 키로 지정하는 것과 같습니다. 기본 키는 파티셔닝에, 기본·보조 키의 조합은 정렬에 사용됩니다. 간단한 예시는 다음과 같습니다.
map의 출력(키):
11.12.1.2
11.14.2.3
11.11.4.1
11.12.1.1
11.14.2.2
3개 리듀서로 파티션(처음 2개 필드가 파티션 키로 사용):
11.11.4.1
-----------
11.12.1.2
11.12.1.1
-----------
11.14.2.3
11.14.2.2
리듀서용 각 파티션 안에서 정렬(모든 4개 필드가 정렬에 사용):
11.11.4.1
-----------
11.12.1.1
11.12.1.2
-----------
11.14.2.2
11.14.2.3
Hadoop Comparator 클래스
Hadoop에는 많은 애플리케이션에 유용한 라이브러리 클래스 KeyFieldBasedComparator가 있습니다. 이 클래스는 Unix/GNU Sort가 제공하는 기능의 부분집합을 제공합니다. 예:
mapred streaming \
-D mapreduce.job.output.key.comparator.class=org.apache.hadoop.mapreduce.lib.partition.KeyFieldBasedComparator \
-D stream.map.output.field.separator=. \
-D stream.num.map.output.key.fields=4 \
-D mapreduce.map.output.key.field.separator=. \
-D mapreduce.partition.keycomparator.options=-k2,2nr \
-D mapreduce.job.reduces=1 \
-input myInputDirs \
-output myOutputDir \
-mapper /bin/cat \
-reducer /bin/cat
위 Map/Reduce 작업의 map 출력 키는 보통 "."으로 구분된 네 필드를 가집니다. 하지만 Map/Reduce 프레임워크는 -D mapreduce.partition.keycomparator.options=-k2,2nr 옵션으로 키의 두 번째 필드로 출력을 정렬합니다. 여기서 -n 은 숫자 정렬을, -r 은 결과를 뒤집도록 지정합니다. 간단한 예시는 아래와 같습니다.
map의 출력(키):
11.12.1.2
11.14.2.3
11.11.4.1
11.12.1.1
11.14.2.2
리듀서용 정렬 출력(정렬에 두 번째 필드 사용):
11.14.2.3
11.14.2.2
11.12.1.2
11.12.1.1
11.11.4.1
Hadoop Aggregate 패키지
Hadoop에는 Aggregate라는 라이브러리 패키지가 있습니다. Aggregate는 특별한 리듀서 클래스와 특별한 컴바이너 클래스, 그리고 값 시퀀스에 대해 "sum", "max", "min" 등의 집계를 수행하는 간단한 집계기 목록을 제공합니다. Aggregate는 매퍼의 각 입력 키/값 쌍에 대해 "aggregatable items"를 생성할 것으로 기대되는 매퍼 플러그인 클래스를 정의하도록 합니다. 컴바이너/리듀서는 적절한 집계기를 호출해 그 집계 가능 항목들을 집계합니다.
Aggregate를 사용하려면 "-reducer aggregate"를 지정하기만 하면 됩니다:
mapred streaming \
-input myInputDirs \
-output myOutputDir \
-mapper myAggregatorForKeyCount.py \
-reducer aggregate \
-file myAggregatorForKeyCount.py
python 프로그램 myAggregatorForKeyCount.py는 다음과 같습니다:
#!/usr/bin/python3
import sys
def generateLongCountToken(id):
return "LongValueSum:" + id + "\t" + "1"
def main(argv):
line = sys.stdin.readline()
try:
while line:
line = line[:-1]
fields = line.split("\t")
print(generateLongCountToken(fields[0]))
line = sys.stdin.readline()
except "end of file":
return None
if __name__ == "__main__":
main(sys.argv)
Hadoop Field Selection 클래스
Hadoop에는 unix "cut" 유틸리티처럼 텍스트 데이터를 처리할 수 있게 해 주는 라이브러리 클래스 FieldSelectionMapReduce가 있습니다. 클래스에 정의된 map 함수는 각 입력 키/값 쌍을 필드 목록으로 취급합니다. 필드 구분자(기본 탭 문자)를 지정할 수 있고, map 출력 키로 임의 필드 목록을, map 출력 값으로 임의 필드 목록을 선택할 수 있습니다. 마찬가지로 클래스에 정의된 reduce 함수는 각 입력 키/값 쌍을 필드 목록으로 취급합니다. reduce 출력 키로 임의 필드 목록을, reduce 출력 값으로 임의 필드 목록을 선택할 수 있습니다. 예:
mapred streaming \
-D mapreduce.map.output.key.field.separator=. \
-D mapreduce.partition.keypartitioner.options=-k1,2 \
-D mapreduce.fieldsel.data.field.separator=. \
-D mapreduce.fieldsel.map.output.key.value.fields.spec=6,5,1-3:0- \
-D mapreduce.fieldsel.reduce.output.key.value.fields.spec=0-2:5- \
-D mapreduce.map.output.key.class=org.apache.hadoop.io.Text \
-D mapreduce.job.reduces=12 \
-input myInputDirs \
-output myOutputDir \
-mapper org.apache.hadoop.mapred.lib.FieldSelectionMapReduce \
-reducer org.apache.hadoop.mapred.lib.FieldSelectionMapReduce \
-partitioner org.apache.hadoop.mapred.lib.KeyFieldBasedPartitioner
"-D mapreduce.fieldsel.map.output.key.value.fields.spec=6,5,1-3:0-" 옵션은 map 출력의 키/값 선택을 지정합니다. 키 선택 스펙과 값 선택 스펙은 ":"으로 구분됩니다. 이 경우 map 출력 키는 필드 6, 5, 1, 2, 3으로 구성됩니다. map 출력 값은 (0-은 필드 0과 이후의 모든 필드) 모든 필드로 구성됩니다.
"-D mapreduce.fieldsel.reduce.output.key.value.fields.spec=0-2:5-" 옵션은 reduce 출력의 키/값 선택을 지정합니다. 이 경우 reduce 출력 키는 필드 0, 1, 2(원래 필드 6, 5, 1에 해당)로 구성됩니다. reduce 출력 값은 필드 5부터 시작하는 모든 필드(모든 원래 필드에 해당)로 구성됩니다.
자주 묻는 질문 (Frequently Asked Questions)
Hadoop Streaming으로 임의의 (준)독립 작업 집합을 실행하려면 어떻게 하나요?
종종 Map Reduce의 완전한 기능이 필요하지 않고, 같은 프로그램의 여러 인스턴스를 — 데이터의 서로 다른 부분에서, 또는 같은 데이터에서 다른 파라미터로 — 실행하기만 하면 됩니다. Hadoop Streaming으로 이렇게 할 수 있습니다.
파일 하나당 map 한 개로 처리하려면 어떻게 하나요?
예로, hadoop 클러스터 전체에 걸쳐 파일 집합을 압축(zipping)하는 문제를 생각해 봅시다. Hadoop Streaming과 사용자 정의 매퍼 스크립트로 이를 달성할 수 있습니다.
- 입력 파일의 전체 HDFS 경로를 담은 파일을 생성합니다. 각 map 작업은 입력으로 파일 이름 하나를 받습니다.
- 파일 이름이 주어지면 그 파일을 로컬 디스크로 가져와 gzip으로 압축하고 원하는 출력 디렉터리에 다시 넣는 매퍼 스크립트를 만듭니다.
리듀서는 몇 개를 써야 하나요?
자세한 내용은 MapReduce Tutorial의 Reducer를 참고하세요.
셸 스크립트에 alias를 설정하면 -mapper 뒤에서 동작할까요?
예: alias c1='cut -f1'을 했다고 합시다. -mapper "c1"이 동작할까요?
alias 사용은 동작하지 않지만, 이 예시에서 보여주듯 변수 치환은 허용됩니다.
$ hdfs dfs -cat /user/me/samples/student_marks
alice 50
bruce 70
charlie 80
dan 75
$ c2='cut -f2'; mapred streaming \
-D mapreduce.job.name='Experiment' \
-input /user/me/samples/student_marks \
-output /user/me/samples/student_out \
-mapper "$c2" -reducer 'cat'
$ hdfs dfs -cat /user/me/samples/student_out/part-00000
50
70
75
80
UNIX 파이프를 쓸 수 있나요?
예: -mapper "cut -f1 | sed s/foo/bar/g"가 동작할까요?
현재 이것은 동작하지 않으며 "java.io.IOException: Broken pipe" 오류를 줍니다. 아마 조사가 필요한 버그입니다.
"No space left on device" 오류가 나면 어떻게 하나요?
예: -file 옵션으로 큰 실행 파일(예: 3.6G)을 배포해 streaming 작업을 실행하면 "No space left on device" 오류가 납니다.
jar 패키징은 구성 변수 stream.tmpdir이 가리키는 디렉터리에서 일어납니다. stream.tmpdir의 기본값은 /tmp입니다. 더 많은 공간이 있는 디렉터리로 값을 설정하세요:
-D stream.tmpdir=/export/bigspace/…
여러 입력 디렉터리를 지정하려면 어떻게 하나요?
여러 '-input' 옵션으로 여러 입력 디렉터리를 지정할 수 있습니다.
mapred streaming \
-input '/user/foo/dir1' -input '/user/foo/dir2' \
(rest of the command)
gzip 형식으로 출력 파일을 생성하려면 어떻게 하나요?
일반 텍스트 파일 대신 gzip 파일을 출력으로 생성할 수 있습니다. streaming 작업 옵션으로 '-D mapreduce.output.fileoutputformat.compress=true -D mapreduce.output.fileoutputformat.compress.codec=org.apache.hadoop.io.compress.GzipCodec'를 전달하세요.
streaming으로 자신의 입출력 포맷을 제공하려면 어떻게 하나요?
패키징하고 사용자 정의 jar를 $HADOOP_CLASSPATH에 넣어 사용자 정의 클래스를 지정할 수 있습니다.
streaming으로 XML 문서를 파싱하려면 어떻게 하나요?
기록 리더 StreamXmlRecordReader를 사용해 XML 문서를 처리할 수 있습니다.
mapred streaming \
-inputreader "StreamXmlRecord,begin=BEGIN_STRING,end=END_STRING" \
(rest of the command)
BEGIN_STRING과 END_STRING 사이에서 찾은 모든 것은 map 작업의 레코드 하나로 처리됩니다.
StreamXmlRecordReader가 이해하는 name-value 속성은 다음과 같습니다.
- (문자열) 'begin' — 레코드 시작을 표시하는 문자, 'end' — 레코드 끝을 표시하는 문자.
- (불리언) 'slowmatch' — 일반 태그 대신 CDATA 안에서 begin과 end 문자를 검색하도록 토글. 기본 false.
- (정수) 'lookahead' — 'slowmatch' 사용 시 CDATA를 동기화할 최대 lookahead 바이트. 'maxrec'보다 커야 함. 기본 2*'maxrec'.
- (정수) 'maxrec' — 'slowmatch' 중 각 일치 사이에 읽을 최대 레코드 크기. 기본 50000바이트.
streaming 애플리케이션에서 카운터를 갱신하려면 어떻게 하나요?
streaming 프로세스는 stderr를 사용해 카운터 정보를 출력할 수 있습니다. reporter:counter:<group>,<counter>,<amount>를 stderr로 보내 카운터를 갱신합니다.
streaming 애플리케이션에서 상태를 갱신하려면 어떻게 하나요?
streaming 프로세스는 stderr를 사용해 상태 정보를 출력할 수 있습니다. 상태를 설정하려면 reporter:status:<message>를 stderr로 보냅니다.
streaming 작업의 매퍼/리듀서에서 Job 변수를 얻으려면 어떻게 하나요?
Configured Parameters를 참고하세요. streaming 작업 실행 중 "mapred" 파라미터의 이름은 변환됩니다. 점( . )이 밑줄( _ )이 됩니다. 예: mapreduce.job.id는 mapreduce_job_id가, mapreduce.job.jar는 mapreduce_job_jar가 됩니다. 코드에서는 밑줄이 있는 파라미터 이름을 사용하세요.
"error=7, Argument list too long" 오류가 나면 어떻게 하나요?
작업은 전체 구성을 환경으로 복사합니다. 작업이 많은 수의 입력 파일을 처리하면 구성 추가가 환경을 넘칠 수 있어요. 환경의 작업 구성 복사본은 작업 실행에 필수적이지 않으며 다음을 설정해 잘라낼 수 있습니다:
-D stream.jobconf.truncate.limit=20000
기본적으로 값은 잘리지 않습니다(-1). 0을 설정하면 이름만 복사하고 값을 복사하지 않습니다. 거의 모든 경우 20000은 환경 넘침을 막는 안전한 값입니다.
더 알아보기 (Learn more)
- 원문: 문서