Apache Spark — 느린 잡의 답은 실행 계획과 이벤트 로그에 있다 · 셔플과 파티션 수 · 실습
셔플 파티션 200개가 만든 파일과 태스크를 줄인다
목표
고객별 매출 집계를 AQE 를 끈 기본값(셔플 파티션 200)·손으로 줄인 값(8)·AQE 로 세 번 돌려, 셔플 뒤 태스크 수와 결과 파일 수가 어떻게 달라지는지 이벤트 로그로 잰다. repartition 과 coalesce 가 셔플 여부와 파일 수에서 어떻게 다른지도 확인한다.
왜 중요한가
groupBy·join 같은 넓은 변환은 같은 키를 한 태스크로 모아야 하므로 셔플을 만든다. 앞 스테이지가 결과를 파티션별 파일로 쓰고, 뒤 스테이지가 그 조각들을 모아 읽는다. 뒤 스테이지의 태스크 수는 spark.sql.shuffle.partitions 가 정하고 기본은 200 이다.
200 은 큰 클러스터를 기준으로 한 숫자다. 몇 MB 자료에 200 을 쓰면 태스크 200개가 각자 조금씩 일하고, 결과를 파일로 쓰면 작은 파일이 200개 생긴다. 반대로 수백 GB 에 200 이면 태스크 하나가 수 GB 를 떠안고 디스크로 흘린다. 그래서 이 숫자는 자료 크기에 맞춰야 하고, AQE 는 셔플이 끝난 뒤 실제 크기를 보고 작은 조각들을 합쳐 이 일을 대신 한다.
파티션 수를 바꾸는 방법은 둘이다. repartition(n) 은 셔플로 조각을 새로 만들고, coalesce(n) 은 셔플 없이 이웃한 조각을 붙인다. coalesce 는 싸지만 줄이기만 할 수 있다 — 조각 둘을 넷으로 만들 수는 없다.
단계
1. /root/spk/shuffle/common.py 에 고객별 매출 DataFrame 을 만드는 함수를 두고, /root/spk/shuffle/noaqe.py(앱 spk-shuffle-noaqe, spark.sql.adaptive.enabled=false)로 결과를 /root/spk/shuffle/out/by_customer_200 에 Parquet 으로 쓰세요.
2. 그 결과 폴더의 데이터 파일(part- 로 시작) 수를 세어 /root/spk/shuffle/out/files_200.txt 에 정수로 적으세요.
3. /root/spk/shuffle/eight.py(앱 spk-shuffle-8, AQE 끔, spark.sql.shuffle.partitions=8)로 같은 결과를 /root/spk/shuffle/out/by_customer_8 에 쓰세요.
4. /root/spk/shuffle/aqe.py(앱 spk-shuffle-aqe, 설정 기본값)로 같은 결과를 /root/spk/shuffle/out/by_customer_aqe 에 쓰고, 그 앱에서 셔플을 읽은 태스크 수를 /root/spk/shuffle/out/aqe_tasks.txt 에 정수로 적으세요.
5. 클릭 원본을 /root/spk/shuffle/repart.py(앱 spk-shuffle-repart)로 repartition(4) 해서 /root/spk/shuffle/out/clicks_repart 에, /root/spk/shuffle/coalesce.py(앱 spk-shuffle-coalesce)로 coalesce(4) 해서 /root/spk/shuffle/out/clicks_coalesce 에 Parquet 으로 쓰고, 두 폴더의 데이터 파일 수를 /root/spk/shuffle/out/repart.json 에 {"repartition_files": 정수, "coalesce_files": 정수} 로 쓰세요.
6. 3단계 앱이 셔플로 쓴 바이트 합을 로그에서 구해 /root/spk/shuffle/out/shuffle_bytes.json 에 {"app": "spk-shuffle-8", "shuffle_write_bytes": 정수} 로 쓰세요.
7. /root/spk/shuffle/narrow.py(앱 spk-shuffle-narrow)로 결제 완료 주문에 day 칼럼을 더해 order_id,customer_id,qty,day 만 /root/spk/shuffle/out/narrow 에 쓰세요. 이 앱에는 셔플이 없어야 합니다.
8. /root/spk/shuffle/report.md 에 ## 200 개의 파티션 ## AQE 가 합친 것 ## repartition 과 coalesce 세 절을 쓰세요. 각 절에 2·4·5단계의 숫자를 넣으세요.
참고
- 스크립트는
/root/spk/shuffle에 두고 거기서spark-submit하세요(from common import …). - 설정은
SparkSession.builder.config("키", "값")이나spark-submit --conf 키=값으로 줍니다. 어느 쪽이든 이벤트 로그에 남습니다. - 로그에서 숫자를 꺼내는 작은 도구
logtool.py를 만들어 두면 4·6단계가 한 줄이 됩니다. 셔플을 읽은 태스크는SparkListenerTaskEnd의Shuffle Read Metrics.Total Records Read가 0 보다 큰 태스크, 셔플로 쓴 바이트는Shuffle Write Metrics.Shuffle Bytes Written의 합입니다. - 흔한 실수: 같은 이름으로 여러 번 돌리고 옛 로그를 세는 것(가장 최근 것을 세세요),
_SUCCESS나.crc까지 파일 수에 넣는 것, coalesce 로 파티션을 늘릴 수 있다고 여기는 것. - 공식 문서: [RDD Programming Guide — Shuffle operations](https://spark.apache.org/docs/4.2.0/rdd-programming-guide.html#shuffle-operations) · [Performance Tuning — Adaptive Query Execution](https://spark.apache.org/docs/4.2.0/sql-performance-tuning.html#adaptive-query-execution) · [Configuration](https://spark.apache.org/docs/4.2.0/configuration.html) · [Monitoring — REST API·metrics](https://spark.apache.org/docs/4.2.0/monitoring.html)
8단계
- AQE 를 끄고 기본값 200 으로
- 200 이 만든 파일 세기
- 손으로 8 로 줄이기
- AQE 가 셔플 뒤에 합친다
- repartition 과 coalesce — 셔플이 있고 없고
- 셔플이 디스크에 쓴 바이트
- 좁은 변환만 있으면 스테이지가 하나다
- 파티션 수를 정한 근거 남기기