Apache Hadoop — Stand up and run HDFS and YARN in one pod
Aggregate access logs with MapReduce written in Python
한국어 원문으로 표시합니다.
목표
Hadoop Streaming 으로 파이썬 맵퍼·리듀서를 써서 접근 로그 사흘치를 상태 코드별로 센다. 컴바이너가 셔플을 얼마나 줄이는지 카운터로 재고, 날짜로 리듀서를 고르는 분할(KeyFieldBasedPartitioner), 깨진 줄을 세는 사용자 카운터, 일부러 죽는 맵퍼로 태스크 실패와 재시도까지 다룬다.
왜 중요한가
Streaming 은 맵과 리듀스를 표준 입출력 계약으로 바꾼다. 맵퍼는 입력 줄을 표준 입력으로 받아 '키값' 줄을 표준 출력으로 내고, 프레임워크가 키로 정렬해 리듀서의 표준 입력에 흘려 준다. 리듀서가 받는 것은 키별 묶음이 아니라 정렬된 줄의 흐름이라, 키가 바뀌는 곳을 스스로 찾아야 한다. 이 계약만 지키면 어떤 언어로도 맵리듀스를 쓸 수 있고, 로컬에서 cat | mapper | sort | reducer 로 똑같이 시험할 수 있다.
컴바이너는 맵 쪽에서 미리 합치는 리듀서다. 결과가 바뀌지 않으려면 연산이 결합·교환 법칙을 만족해야 하고(합·최댓값은 되고 평균은 안 된다), 입력과 출력의 모양이 같아야 한다. 맞게 붙이면 셔플이 몇 십 분의 일로 준다.
기본 분할은 키 전체의 해시다. 키가 '날짜경로' 인데 날짜 하나의 모든 경로를 한 리듀서에서 보고 싶다면, 키의 첫 칸으로만 분할하도록 바꿔야 한다. 그리고 운영에서 스트리밍 잡이 조용히 틀리는 가장 흔한 이유는 깨진 줄이다 — 버리는 것은 괜찮지만 몇 줄을 버렸는지 세지 않는 것은 괜찮지 않다.
단계
- /root/hdp/streaming/mapper.py 를 쓰세요. 한 줄을 공백으로 나눠 칸이 10개 이상이고 아홉 번째 칸(0부터 8번)이 세 자리 숫자이며
#으로 시작하지 않는 줄이면상태코드<TAB>1을 내고, 아니면 버립니다. - /root/hdp/streaming/reducer.py 를 쓰세요. 키로 정렬된
키<TAB>수줄을 받아 같은 키의 수를 더해키<TAB>합을 냅니다(수가 1 이 아니어도 됩니다). /data/logs/access-2026-03-01.log~03.log를 HDFS /user/root/streaming/in/ 에 올리고, 잡 이름hdp-status로 스트리밍 잡을 돌려 /user/root/streaming/status 에 쓰세요.- 같은 잡에
-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": …}로 쓰세요. - /root/hdp/streaming/mapper2.py 가
날짜(yyyy-MM-dd)<TAB>경로<TAB>1을 내게 하고, 1단계의 reducer.py 를 그대로 리듀서로 써서 잡 이름hdp-daily, 리듀서 2개, 키 두 칸, 첫 칸으로만 분할하도록 돌려 /user/root/streaming/daily 에 쓰세요. - 1단계 맵퍼에 깨진 줄마다 사용자 카운터
Lab·BadLines를 1 올리는 줄을 넣고(표준 오류에reporter:counter:Lab,BadLines,1), 잡 이름hdp-bad로 /user/root/streaming/bad 에 돌리세요. /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 에 다시 돌려 성공시키세요.- /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 · MapReduce Tutorial · mapred-default.xml
맵퍼 — 한 줄 받아 키와 값을 내기
/root/hdp/streaming/mapper.py 를 쓰세요. 표준 입력의 줄마다 공백으로 나눠, 칸이 10개 이상이고 9번째 칸(0부터 세어 8)이 세 자리 숫자이며 # 으로 시작하지 않으면 상태코드<TAB>1 을 표준 출력으로 내고, 아니면 아무것도 내지 않습니다.
head -1000 /data/logs/access-2026-03-01.log | python3 mapper.py | sort | uniq -c 로 먼저 시험하세요. 채점기는 여러분의 맵퍼를 깨진 줄이 섞인 새 입력에 직접 돌려 봅니다.
리듀서 — 정렬된 흐름에서 키가 바뀌는 곳 찾기
/root/hdp/streaming/reducer.py 를 쓰세요. 표준 입력으로 키 순으로 정렬된 키<TAB>수 줄을 받아, 같은 키의 수를 더해 키<TAB>합 을 표준 출력으로 냅니다. 수는 1 이 아닐 수도 있습니다.
리듀서는 키별 목록을 받지 않습니다. 정렬된 줄이 흘러들어올 뿐이라, 앞 줄과 키가 달라지는 순간 앞 키의 합을 내보내고 새로 시작합니다. 마지막 키를 내보내는 것을 잊지 마세요. 수가 1 이 아니어도 되게 쓰면 그대로 컴바이너가 됩니다.
YARN 에 올리기
lab-hadoop start yarn 으로 YARN 을 켜고, /data/logs/access-2026-03-01.log~03.log 를 HDFS /user/root/streaming/in/ 에 올린 뒤, mapred streaming -D mapreduce.job.name=hdp-status -files mapper.py,reducer.py -mapper 'python3 mapper.py' -reducer 'python3 reducer.py' -input /user/root/streaming/in -output /user/root/streaming/status 로 돌리세요.
-files 로 준 스크립트는 각 태스크의 작업 폴더에 복사됩니다. 그래서 -mapper 에는 경로 없이 파일 이름만 씁니다. 채점기는 출력이 원본 로그에서 직접 센 값과 같은지, 이력에 성공한 hdp-status 잡이 있는지 봅니다.
컴바이너로 셔플 줄이기
3단계 잡에 -combiner 'python3 reducer.py' 를 더해 잡 이름 hdp-status-comb, 출력 /user/root/streaming/status_comb 로 돌리세요. 그 잡의 이력에서 TaskCounter 의 MAP_OUTPUT_RECORDS·COMBINE_OUTPUT_RECORDS·REDUCE_INPUT_RECORDS 를 꺼내 /root/hdp/streaming/shuffle.json 에 {"map_output": 정수, "combine_output": 정수, "reduce_input": 정수} 로 쓰세요.
합은 결합·교환 법칙을 만족하므로 리듀서를 그대로 컴바이너로 써도 결과가 같습니다. 컴바이너 뒤에 셔플로 넘어간 줄 수(=리듀스 입력)가 맵 출력보다 얼마나 적은지 보세요.
날짜로 리듀서 고르기
/root/hdp/streaming/mapper2.py 가 유효한 줄마다 날짜(yyyy-MM-dd)<TAB>경로<TAB>1 을 내게 하고, reducer.py 를 리듀서로 써서 잡 이름 hdp-daily, -D stream.num.map.output.key.fields=2 -D mapreduce.partition.keypartitioner.options=-k1,1(일반 옵션, 앞쪽)과 -partitioner org.apache.hadoop.mapred.lib.KeyFieldBasedPartitioner -numReduceTasks 2(명령 옵션, 뒤쪽)로 /user/root/streaming/daily 에 쓰세요.
키는 앞 두 칸(날짜·경로)이라 정렬은 두 칸으로 하지만, 분할은 첫 칸(날짜)으로만 합니다. 그래서 한 날짜의 모든 경로가 한 리듀서로 가고, 출력 파일 하나가 날짜 하나를 통째로 가집니다. reducer.py 가 마지막 탭에서 나누도록 썼다면 키가 두 칸이어도 그대로 동작합니다.
깨진 줄을 사용자 카운터로 세기
mapper.py 가 깨진 줄마다 표준 오류로 reporter:counter:Lab,BadLines,1 을 쓰게 하고, 잡 이름 hdp-bad 로 /user/root/streaming/bad 에 돌리세요. 그 잡의 카운터 Lab·BadLines 가 세 로그의 깨진 줄 수와 같아야 합니다.
스트리밍은 표준 오류의 reporter:counter:<그룹>,<이름>,<증가량> 줄을 카운터 갱신으로 알아듣습니다. 이 카운터는 잡 이력에 남아, 잡이 끝난 뒤에도 '몇 줄을 버렸는가' 에 답할 수 있습니다.
죽는 맵퍼 — 실패와 다시 돌리기
/event/spring 경로를 만나면 예외로 끝나는 /root/hdp/streaming/crash_mapper.py 를 써서 잡 이름 hdp-crash, -D mapreduce.map.maxattempts=1 로 /user/root/streaming/crash 에 돌려 실패시키고, 그 잡 ID 를 /root/hdp/streaming/crash.txt 에 적으세요. 그다음 같은 잡 이름으로 1단계의 mapper.py 를 써서 /user/root/streaming/crash-ok 에 다시 돌려 성공시키세요.
맵 태스크가 실패하면 프레임워크는 다른 시도로 다시 돌립니다(기본 4번). 시도를 1번으로 줄이면 첫 실패가 곧 잡 실패입니다. 실패한 잡도 이력이 남으니, 무엇이 왜 죽었는지는 yarn logs -applicationId application_<같은 숫자> 로 태스크의 표준 오류를 읽어 확인합니다.
계약·컴바이너·실패를 남기기
/root/hdp/streaming/report.md 에 ## 표준 입출력 계약 ## 컴바이너와 분할 ## 실패와 카운터 세 절을 쓰세요. 둘째 절에 4단계의 map_output 과 reduce_input 을 숫자로 넣으세요.
리듀서가 받는 것이 무엇인지, 컴바이너가 셔플을 몇 분의 일로 줄였는지, 깨진 줄과 실패한 태스크를 어떻게 알아챘는지를 적으세요.