Apache Spark — The answer to a slow job is in the plan and the event log
Launch a first job and check it in the event log — transformations are lazy, actions do the work
한국어 원문으로 표시합니다.
목표
local 모드 Spark 로 주문 30만 건을 읽어 세고, 변환만 쌓은 앱과 행동을 부른 앱을 이벤트 로그로 견준다. 잡·스테이지·태스크가 언제 생기고 파티션 수가 무엇으로 정해지는지 숫자로 확인한다.
왜 중요한가
Spark 코드를 처음 읽으면 한 줄 한 줄이 그 자리에서 실행되는 것처럼 보인다. 그렇지 않다. filter·withColumn·select 같은 변환은 계획에 한 줄을 더할 뿐이고, count·take·write 같은 행동을 부를 때 그 계획 전체가 잡으로 바뀌어 실행된다. 이것을 모르면 "읽기가 느리다" 고 보이는 곳이 사실은 앞에서 쌓아 둔 변환 전체가 한꺼번에 도는 자리라는 것을 놓친다.
그 차이는 눈으로 확인할 수 있다. Spark 는 앱이 끝나면 무슨 일을 했는지를 이벤트 로그로 남긴다. 잡이 몇 개 떴는지, 스테이지가 몇 개였는지, 태스크가 몇 개였는지, 어떤 물리 계획을 썼는지가 전부 JSON 한 줄씩 들어 있다. 운영에서 끝난 잡을 되짚을 때 보는 것도 이것이다(히스토리 서버가 이 파일을 읽는다).
local 모드에서는 드라이버 JVM 하나가 곧 실행기다. local[2] 는 코어 둘을 쓴다는 뜻이고, 이 숫자가 기본 병렬도와 파일을 몇 조각으로 나눌지를 함께 정한다.
단계
spark-submit --version의 출력(표준 오류 포함)을 /root/spk/first/version.txt 에 저장하세요.- /root/spk/first/count.py 를 만들어 앱 이름
spk-first-count로/data/shop/orders.csv(머리줄 있음)의 행 수를 세고 /root/spk/first/out/count.json 에{"rows": 정수}로 쓰세요. - /root/spk/first/lazy.py 를 만들어 앱 이름
spk-first-lazy로 스키마를 직접 주어 읽고filter(또는where)와withColumn을 쌓되 행동은 부르지 마세요. 실행한 뒤 그 앱의 잡이 0개여야 합니다. - /root/spk/first/actions.py 를 만들어 앱 이름
spk-first-actions로 행동을 세 번 이상(예:count·take·write) 부르고, 100줄을 /root/spk/first/out/sample 에 JSON 으로 쓰세요. 그 앱의 이벤트 로그에서 잡 시작 이벤트 수를 세어 /root/spk/first/out/jobs.txt 에 정수 하나로 적으세요. - /root/spk/first/agg.py 를 만들어 앱 이름
spk-first-agg로status별 주문 수를 /root/spk/first/out/by_status 에 머리줄 있는 CSV 로 쓰세요. 그 앱에서 완료된 스테이지 수를 /root/spk/first/out/stages.txt 에 정수로 적으세요. - /root/spk/first/par.py 를 만들어 첫 인자를 앱 이름으로 받게 하고,
--master local[1]로spk-first-par1,--master local[2]로spk-first-par2를 돌려 각각 /root/spk/first/out/spk-first-par1.json·/root/spk/first/out/spk-first-par2.json 에master·default_parallelism·partitions(주문 파일을 읽은 DataFrame 의 파티션 수)를 쓰세요. - /root/spk/first/fail.py 를 만들어 앱 이름
spk-first-fail로 없는 경로/data/shop/order.csv를 읽게 하고, 잡은 예외의 오류 조건 이름(getCondition())을 /root/spk/first/out/error.txt 첫 줄에 쓰세요. - /root/spk/first/report.md 에
## 잡과 스테이지## 게으른 실행## 파티션세 절로 여러분의 숫자를 적으세요. 첫 절에는 4단계의 잡 수를, 셋째 절에는 6단계의 파티션 수 둘을 숫자로 넣으세요.
참고
- 실행은
spark-submit /root/spk/first/count.py처럼 합니다. 기본 설정은/opt/spark/conf/spark-defaults.conf에 있습니다(local[2], 드라이버 메모리 1g, 이벤트 로그 켜짐). - 이벤트 로그는 앱이 끝나면
/root/spark-events/local-<숫자>로 남습니다. 도는 중에는 이름 끝에.inprogress가 붙습니다. 어느 파일이 어느 앱인지는grep -l '"App Name":"spk-first-actions"' /root/spark-events/*로 찾습니다. - 잡 수는
grep -c '"Event":"SparkListenerJobStart"' <파일>, 스테이지 수는"Event":"SparkListenerStageCompleted"를 셉니다.jq -c 'select(.Event=="SparkListenerJobStart")'로 내용을 볼 수도 있습니다. - 흔한 실수:
spark.stop()을 빼서 로그가.inprogress로 남는 것, 머리줄 있는 CSV 에서header=True를 빼서 머리줄을 한 행으로 세는 것, 스키마 없이 읽어 머리줄을 확인하는 잡이 생기는 것(3단계는 잡이 0개여야 합니다). - 공식 문서: Cluster Mode Overview · RDD Programming Guide — RDD Operations · Monitoring — Viewing After the Fact · Submitting Applications — Master URLs
어떤 Spark 인지 확인하기
spark-submit --version 의 출력을 표준 오류까지 합쳐 /root/spk/first/version.txt 에 저장하세요.
Spark 는 판 정보를 표준 오류로 찍습니다. 2>&1 로 합치지 않으면 파일이 비어 버립니다. 판 번호와 함께 어떤 Scala·자바 위에서 도는지도 보입니다 — Spark 4 는 자바 17 이상을 요구합니다.
첫 잡 — 30만 건 세기
/root/spk/first/count.py 를 앱 이름 spk-first-count 로 만들어 /data/shop/orders.csv(머리줄 있음)의 행 수를 세고 /root/spk/first/out/count.json 에 {"rows": 정수} 로 쓰세요. spark-submit 으로 실행하세요.
SparkSession.builder.appName(...) 이 앱 이름을 정합니다. 채점기는 그 이름의 이벤트 로그가 끝까지 닫혔는지와 파일의 숫자가 원본의 행 수와 같은지를 봅니다. 머리줄을 행으로 세면 하나가 많습니다.
변환만 쌓으면 잡이 뜨지 않는다
/root/spk/first/lazy.py 를 앱 이름 spk-first-lazy 로 만들어 스키마를 직접 준 spark.read.csv 로 주문을 읽고, filter(또는 where)와 withColumn 을 쌓되 행동은 하나도 부르지 마세요. 실행한 뒤 그 앱의 이벤트 로그에 잡이 0개여야 합니다.
스키마를 주지 않으면 Spark 가 머리줄을 확인하려고 작은 잡을 띄웁니다. schema="order_id STRING, ..." 처럼 DDL 문자열을 넘기면 그 잡도 사라집니다. 스키마를 출력하는 것은 괜찮습니다 — 스키마는 드라이버가 이미 알고 있어 자료를 읽지 않습니다.
행동마다 잡이 뜬다 — 이벤트 로그로 세기
/root/spk/first/actions.py 를 앱 이름 spk-first-actions 로 만들어 행동을 세 번 이상 부르고, 그중 하나로 100줄을 /root/spk/first/out/sample 에 JSON 으로 쓰세요. 실행한 뒤 그 앱의 이벤트 로그에서 SparkListenerJobStart 이벤트 수를 세어 /root/spk/first/out/jobs.txt 에 정수 하나로 적으세요.
행동 하나가 잡 하나라고 단정하지 마세요. 파일을 쓰는 행동이나 스키마 추론처럼 잡을 둘 이상 만드는 행동도 있습니다. 그래서 세는 것은 코드가 아니라 로그입니다. 같은 이름으로 여러 번 돌렸다면 가장 최근 로그를 세야 합니다.
넓은 변환이 스테이지를 가른다
/root/spk/first/agg.py 를 앱 이름 spk-first-agg 로 만들어 status 별 주문 수를 /root/spk/first/out/by_status 에 머리줄 있는 CSV(칸: status,count)로 쓰세요. 그 앱의 이벤트 로그에서 SparkListenerStageCompleted 이벤트 수를 세어 /root/spk/first/out/stages.txt 에 정수로 적으세요.
groupBy 는 같은 키를 한 태스크로 모아야 하므로 셔플이 생기고, 셔플 앞뒤가 서로 다른 스테이지가 됩니다. 채점기는 여러분의 CSV 가 원본에서 센 값과 같은지, 그 앱에 셔플을 쓴 스테이지가 실제로 있었는지, 적은 스테이지 수가 로그와 같은지를 봅니다.
코어 수가 파티션 수를 정한다
/root/spk/first/par.py 를 첫 인자를 앱 이름으로 받게 만들고, spark-submit --master 'local[1]' par.py spk-first-par1 과 spark-submit --master 'local[2]' par.py spk-first-par2 로 두 번 돌려 /root/spk/first/out/spk-first-par1.json·/root/spk/first/out/spk-first-par2.json 에 {"master": 문자열, "default_parallelism": 정수, "partitions": 정수} 를 쓰세요. partitions 는 스키마를 주어 읽은 주문 DataFrame 의 rdd.getNumPartitions() 입니다.
파일을 몇 조각으로 나눌지는 spark.sql.files.maxPartitionBytes(기본 128MB)만으로 정해지지 않습니다. 코어가 여럿이면 전체 크기를 코어 수로 나눈 값이 더 작을 때 그 크기로 자릅니다 — 코어를 놀리지 않으려는 규칙입니다. 16MB 파일이 코어 하나와 둘에서 어떻게 달라지는지 보세요.
실패도 앱이다 — 오류 조건 이름 읽기
/root/spk/first/fail.py 를 앱 이름 spk-first-fail 로 만들어 없는 경로 /data/shop/order.csv 를 읽게 하고, 예외를 잡아 getCondition() 이 돌려준 오류 조건 이름을 /root/spk/first/out/error.txt 첫 줄에 쓰세요. 앱은 spark.stop() 으로 정상 종료해야 합니다.
없는 경로는 행동을 부를 때가 아니라 read 하는 순간 드러납니다. 파일 목록은 계획을 세울 때 확인하기 때문입니다. Spark 4 의 오류는 사람이 읽는 문장 말고도 기계가 읽는 이름(오류 조건)을 가집니다. 경보나 재시도 규칙은 문장이 아니라 이 이름으로 거는 것이 안전합니다.
무엇을 봤는지 숫자로 남기기
/root/spk/first/report.md 에 ## 잡과 스테이지 ## 게으른 실행 ## 파티션 세 절을 쓰세요. 첫 절에는 4단계에서 센 잡 수를, 셋째 절에는 6단계의 두 파티션 수를 숫자로 넣으세요.
숫자만 옮기지 말고 왜 그 숫자인지 한두 문장씩 붙이세요. 행동 셋에 잡이 몇 개였는지, 변환만 쌓은 앱은 왜 0개였는지, 코어 수가 파티션 수를 어떻게 바꿨는지가 이 실습에서 본 전부입니다.