Apache Hadoop — HDFS 와 YARN 을 한 파드에 세우고 운영한다 · Hadoop Streaming · 实验
접근 로그를 파이썬 맵리듀스로 집계한다
목표
Hadoop Streaming 으로 파이썬 맵퍼·리듀서를 써서 접근 로그 사흘치를 상태 코드별로 센다. 컴바이너가 셔플을 얼마나 줄이는지 카운터로 재고, 날짜로 리듀서를 고르는 분할(KeyFieldBasedPartitioner), 깨진 줄을 세는 사용자 카운터, 일부러 죽는 맵퍼로 태스크 실패와 재시도까지 다룬다.
왜 중요한가
Streaming 은 맵과 리듀스를 표준 입출력 계약으로 바꾼다. 맵퍼는 입력 줄을 표준 입력으로 받아 '키<TAB>값' 줄을 표준 출력으로 내고, 프레임워크가 키로 정렬해 리듀서의 표준 입력에 흘려 준다. 리듀서가 받는 것은 키별 묶음이 아니라 정렬된 줄의 흐름이라, 키가 바뀌는 곳을 스스로 찾아야 한다. 이 계약만 지키면 어떤 언어로도 맵리듀스를 쓸 수 있고, 로컬에서 cat | mapper | sort | reducer 로 똑같이 시험할 수 있다.
컴바이너는 맵 쪽에서 미리 합치는 리듀서다. 결과가 바뀌지 않으려면 연산이 결합·교환 법칙을 만족해야 하고(합·최댓값은 되고 평균은 안 된다), 입력과 출력의 모양이 같아야 한다. 맞게 붙이면 셔플이 몇 십 분의 일로 준다.
기본 분할은 키 전체의 해시다. 키가 '날짜<TAB>경로' 인데 날짜 하나의 모든 경로를 한 리듀서에서 보고 싶다면, 키의 첫 칸으로만 분할하도록 바꿔야 한다. 그리고 운영에서 스트리밍 잡이 조용히 틀리는 가장 흔한 이유는 깨진 줄이다 — 버리는 것은 괜찮지만 몇 줄을 버렸는지 세지 않는 것은 괜찮지 않다.
단계
1. /root/hdp/streaming/mapper.py 를 쓰세요. 한 줄을 공백으로 나눠 칸이 10개 이상이고 아홉 번째 칸(0부터 8번)이 세 자리 숫자이며 # 으로 시작하지 않는 줄이면 상태코드<TAB>1 을 내고, 아니면 버립니다.
2. /root/hdp/streaming/reducer.py 를 쓰세요. 키로 정렬된 키<TAB>수 줄을 받아 같은 키의 수를 더해 키<TAB>합 을 냅니다(수가 1 이 아니어도 됩니다).
3. /data/logs/access-2026-03-01.log~03.log 를 HDFS /user/root/streaming/in/ 에 올리고, 잡 이름 hdp-status 로 스트리밍 잡을 돌려 /user/root/streaming/status 에 쓰세요.
4. 같은 잡에 -combiner 로 리듀서를 붙여 잡 이름 hdp-status-comb, 출력 /user/root/streaming/status_comb 로 돌리고, 그 잡의 MAP_OUTPUT_RECORDS·COMBINE_OUTPUT_RECORDS·REDUCE_INPUT_RECORDS 를 /root/hdp/streaming/shuffle.json 에 {"map_output": …, "combine_output": …, "reduce_input": …} 로 쓰세요.
5. /root/hdp/streaming/mapper2.py 가 날짜(yyyy-MM-dd)<TAB>경로<TAB>1 을 내게 하고, 1단계의 reducer.py 를 그대로 리듀서로 써서 잡 이름 hdp-daily, 리듀서 2개, 키 두 칸, 첫 칸으로만 분할하도록 돌려 /user/root/streaming/daily 에 쓰세요.
6. 1단계 맵퍼에 깨진 줄마다 사용자 카운터 Lab·BadLines 를 1 올리는 줄을 넣고(표준 오류에 reporter:counter:Lab,BadLines,1), 잡 이름 hdp-bad 로 /user/root/streaming/bad 에 돌리세요.
7. /event/spring 경로를 만나면 예외로 끝나는 /root/hdp/streaming/crash_mapper.py 로 잡 이름 hdp-crash, 맵 시도 1회(-D mapreduce.map.maxattempts=1)로 /user/root/streaming/crash 에 돌려 실패시키고, 그 잡 ID 를 /root/hdp/streaming/crash.txt 에 적은 뒤, 같은 이름으로 1단계 맵퍼를 써서 /user/root/streaming/crash-ok 에 다시 돌려 성공시키세요.
8. /root/hdp/streaming/report.md 에 ## 표준 입출력 계약 ## 컴바이너와 분할 ## 실패와 카운터 세 절을 쓰세요. 둘째 절에 4단계의 map_output 과 reduce_input 을 넣으세요.
참고
- 스트리밍 잡:
mapred streaming -D mapreduce.job.name=<이름> -files mapper.py,reducer.py -mapper 'python3 mapper.py' -reducer 'python3 reducer.py' -input <입력> -output <출력>.-D옵션은 다른 옵션보다 앞에 둡니다. - 로컬 시험:
head -1000 /data/logs/access-2026-03-01.log | python3 mapper.py | sort | python3 reducer.py. - 키 두 칸·첫 칸으로 분할: 일반 옵션
-D stream.num.map.output.key.fields=2 -D mapreduce.partition.keypartitioner.options=-k1,1은-files앞에, 명령 옵션-partitioner org.apache.hadoop.mapred.lib.KeyFieldBasedPartitioner -numReduceTasks 2는-mapper뒤에 둡니다. 순서가 틀리면 'Unrecognized option' 으로 멈춥니다. - 잡 이력은
/tmp/hadoop-yarn/staging/history/done_intermediate/root/에 남습니다(잡 이름의-는 파일 이름에서%2D로 적힙니다). 실패한 잡도 남습니다. - 흔한 실수: 리듀서가 키마다 한 번만 불린다고 여기는 것(줄마다 불리는 흐름입니다), 컴바이너 출력 모양을 입력과 다르게 만드는 것, 출력 디렉터리를 지우지 않고 다시 돌리는 것.
- 공식 문서: [Hadoop Streaming](https://hadoop.apache.org/docs/r3.5.0/hadoop-streaming/HadoopStreaming.html) · [MapReduce Tutorial](https://hadoop.apache.org/docs/r3.5.0/hadoop-mapreduce-client/hadoop-mapreduce-client-core/MapReduceTutorial.html) · [mapred-default.xml](https://hadoop.apache.org/docs/r3.5.0/hadoop-mapreduce-client/hadoop-mapreduce-client-core/mapred-default.xml)
8个步骤
- 맵퍼 — 한 줄 받아 키와 값을 내기
- 리듀서 — 정렬된 흐름에서 키가 바뀌는 곳 찾기
- YARN 에 올리기
- 컴바이너로 셔플 줄이기
- 날짜로 리듀서 고르기
- 깨진 줄을 사용자 카운터로 세기
- 죽는 맵퍼 — 실패와 다시 돌리기
- 계약·컴바이너·실패를 남기기