Apache Spark — 느린 잡의 답은 실행 계획과 이벤트 로그에 있다 · 실행 계획 읽기 · 실습
explain 으로 Parquet 을 얼마나 덜 읽는지 확인한다
목표
주문을 Parquet 으로 바꾼 뒤, 같은 질문이 조건을 어떻게 쓰느냐에 따라 파일을 얼마나 덜 읽는지(밀어내기·칼럼 가지치기·파티션 가지치기)를 실행 계획으로 읽는다. 적응형 실행(AQE)이 돌고 난 뒤 계획을 어떻게 바꿨는지도 이벤트 로그로 확인한다.
왜 중요한가
Spark 가 느릴 때 가장 먼저 볼 것은 코드가 아니라 계획이다. 같은 결과를 내는 두 줄이 한쪽은 파일에서 필요한 칼럼과 행만 읽고, 다른 쪽은 전부 읽은 뒤 버린다. 그 차이는 결과에 드러나지 않고 계획에만 드러난다.
Parquet 같은 열 형식은 두 가지를 할 수 있다. 필요한 칼럼만 읽고(칼럼 가지치기, 계획의 ReadSchema), 행 그룹의 최솟값·최댓값 통계로 조건에 안 맞는 덩어리를 건너뛴다(조건 밀어내기, 계획의 PushedFilters). 여기에 디렉터리로 나눈 파티션을 조건으로 고르면 디렉터리째 건너뛴다(PartitionFilters). 셋 다 옵티마이저가 조건을 이해할 수 있을 때만 일어난다. 칼럼을 함수로 감싸는 순간 밀어내기는 사라진다.
AQE 는 셔플이 끝난 뒤 실제 크기를 보고 남은 계획을 고친다. 그래서 돌기 전의 계획(isFinalPlan=false)과 돈 뒤의 계획은 다를 수 있고, 진짜로 무엇이 돌았는지는 돈 뒤의 계획이나 이벤트 로그에서 본다.
단계
1. /root/spk/plan/convert.py(앱 spk-plan-convert)로 주문을 스키마로 읽어 order_date(날짜) 칼럼을 더하고 /root/spk/plan/orders_pq 에 Parquet 으로 쓰세요.
2. 도우미 /root/spk/plan/planhelp.py 에 explain·field·items 함수를 두고, /root/spk/plan/explain.py(앱 spk-plan-explain)로 환불 주문의 order_id·qty 를 고르는 질의의 extended·formatted 계획을 /root/spk/plan/out/extended.txt·/root/spk/plan/out/formatted.txt 에 저장하고 그 질의를 collect() 하세요.
3. /root/spk/plan/pushdown.py(앱 spk-plan-push)로 status == 'refunded' 이고 qty >= 5 인 주문의 order_id,qty,channel 을 /root/spk/plan/out/pushdown 에 CSV 로 쓰고, 그 계획의 PushedFilters 항목 목록과 ReadSchema 를 /root/spk/plan/out/pushdown.json 에 쓰세요.
4. /root/spk/plan/nopush.py(앱 spk-plan-nopush)로 같은 질문을 upper(status) == 'REFUNDED' 로 바꿔 /root/spk/plan/out/nopush 에 쓰고, 그 계획의 PushedFilters 목록을 /root/spk/plan/out/nopush.json 에 쓰세요.
5. /root/spk/plan/partition.py(앱 spk-plan-part)로 /root/spk/plan/orders_by_day 에 order_date 로 나눠 쓰고, 2026-02-14 하루를 골라 센 행 수와 계획의 PartitionFilters 를 /root/spk/plan/out/partition.json 에 쓰세요.
6. /root/spk/plan/aqe.py(앱 spk-plan-aqe)로 고객별 주문 수를 collect() 한 뒤 계획을 /root/spk/plan/out/aqe_final.txt 에 저장하고, 그 앱의 로그에서 셔플을 읽은 태스크 수를 세어 /root/spk/plan/out/aqe.json 에 {"shuffle_partitions": 200, "reduce_tasks": 정수} 로 쓰세요.
7. /root/spk/plan/optimize.py(앱 spk-plan-opt)로 qty > 1 과 qty > 3 두 where 를 잇고 1 + 2 를 three 로 고르는 질의의 extended 계획을 /root/spk/plan/out/optimized.txt 에 저장하고 collect() 하세요.
8. /root/spk/plan/report.md 에 ## 밀어내기 ## 가지치기 ## AQE 세 절을 쓰세요. 둘째 절에는 5단계의 행 수를, 셋째 절에는 6단계의 태스크 수를 숫자로 넣으세요.
참고
- 스크립트는
/root/spk/plan에 두고 거기서spark-submit하세요. 같은 폴더의planhelp.py를from planhelp import explain으로 씁니다. explain의 모드:simple(물리 계획만),extended(파싱·분석·최적화된 논리 계획 + 물리 계획),codegen,cost,formatted(물리 계획 트리 + 연산자별 상세). 연산자별PushedFilters·ReadSchema·PartitionFilters는 formatted 에서 한 줄씩 보입니다.- 이벤트 로그에서 셔플을 읽은 태스크는
SparkListenerTaskEnd의Task Metrics→Shuffle Read Metrics→Total Records Read가 0 보다 큰 태스크입니다. - 흔한 실수: CSV 에서 밀어내기를 찾는 것(CSV 는 파일을 다 읽어야 합니다), 계획을
collect()전에 찍고 AQE 최종 계획이라고 여기는 것,PushedFilters가 있으니 파일을 덜 읽었다고 단정하는 것(행 그룹 통계가 조건을 가를 수 있어야 건너뜁니다). - 공식 문서: [Performance Tuning](https://spark.apache.org/docs/4.2.0/sql-performance-tuning.html) · [EXPLAIN](https://spark.apache.org/docs/4.2.0/sql-ref-syntax-qry-explain.html) · [Parquet Files](https://spark.apache.org/docs/4.2.0/sql-data-sources-parquet.html) · [Adaptive Query Execution](https://spark.apache.org/docs/4.2.0/sql-performance-tuning.html#adaptive-query-execution)
8단계
- CSV 를 Parquet 으로
- explain 의 모드 — 논리 계획에서 물리 계획까지
- 조건 밀어내기와 칼럼 가지치기
- 함수로 감싸면 밀어내기가 사라진다
- 파티션 가지치기 — 디렉터리째 건너뛰기
- AQE 가 돈 뒤에 고친 계획
- 옵티마이저가 합치고 접는 것
- 덜 읽은 만큼을 숫자로