LabHub
배우기 러닝패스 코스

Apache Hadoop — 1つのPodにHDFSとYARNを立てて運用する

アクセスログを Python の MapReduce で集計する

LabHub 에서 이어서 보기

한국어 원문으로 표시합니다.

목표

Hadoop Streaming 으로 파이썬 맵퍼·리듀서를 써서 접근 로그 사흘치를 상태 코드별로 센다. 컴바이너가 셔플을 얼마나 줄이는지 카운터로 재고, 날짜로 리듀서를 고르는 분할(KeyFieldBasedPartitioner), 깨진 줄을 세는 사용자 카운터, 일부러 죽는 맵퍼로 태스크 실패와 재시도까지 다룬다.

왜 중요한가

Streaming 은 맵과 리듀스를 표준 입출력 계약으로 바꾼다. 맵퍼는 입력 줄을 표준 입력으로 받아 '키값' 줄을 표준 출력으로 내고, 프레임워크가 키로 정렬해 리듀서의 표준 입력에 흘려 준다. 리듀서가 받는 것은 키별 묶음이 아니라 정렬된 줄의 흐름이라, 키가 바뀌는 곳을 스스로 찾아야 한다. 이 계약만 지키면 어떤 언어로도 맵리듀스를 쓸 수 있고, 로컬에서 cat | mapper | sort | reducer 로 똑같이 시험할 수 있다. 컴바이너는 맵 쪽에서 미리 합치는 리듀서다. 결과가 바뀌지 않으려면 연산이 결합·교환 법칙을 만족해야 하고(합·최댓값은 되고 평균은 안 된다), 입력과 출력의 모양이 같아야 한다. 맞게 붙이면 셔플이 몇 십 분의 일로 준다. 기본 분할은 키 전체의 해시다. 키가 '날짜경로' 인데 날짜 하나의 모든 경로를 한 리듀서에서 보고 싶다면, 키의 첫 칸으로만 분할하도록 바꿔야 한다. 그리고 운영에서 스트리밍 잡이 조용히 틀리는 가장 흔한 이유는 깨진 줄이다 — 버리는 것은 괜찮지만 몇 줄을 버렸는지 세지 않는 것은 괜찮지 않다.

단계

  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_outputreduce_input 을 넣으세요.

참고

맵퍼 — 한 줄 받아 키와 값을 내기

/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_outputreduce_input 을 숫자로 넣으세요.

리듀서가 받는 것이 무엇인지, 컴바이너가 셔플을 몇 분의 일로 줄였는지, 깨진 줄과 실패한 태스크를 어떻게 알아챘는지를 적으세요.