Apache Spark — 遅いジョブの答えは実行計画とイベントログにある
explain で Parquet をどれだけ読まずに済むか確かめる
한국어 원문으로 표시합니다.
목표
주문을 Parquet 으로 바꾼 뒤, 같은 질문이 조건을 어떻게 쓰느냐에 따라 파일을 얼마나 덜 읽는지(밀어내기·칼럼 가지치기·파티션 가지치기)를 실행 계획으로 읽는다. 적응형 실행(AQE)이 돌고 난 뒤 계획을 어떻게 바꿨는지도 이벤트 로그로 확인한다.
왜 중요한가
Spark 가 느릴 때 가장 먼저 볼 것은 코드가 아니라 계획이다. 같은 결과를 내는 두 줄이 한쪽은 파일에서 필요한 칼럼과 행만 읽고, 다른 쪽은 전부 읽은 뒤 버린다. 그 차이는 결과에 드러나지 않고 계획에만 드러난다.
Parquet 같은 열 형식은 두 가지를 할 수 있다. 필요한 칼럼만 읽고(칼럼 가지치기, 계획의 ReadSchema), 행 그룹의 최솟값·최댓값 통계로 조건에 안 맞는 덩어리를 건너뛴다(조건 밀어내기, 계획의 PushedFilters). 여기에 디렉터리로 나눈 파티션을 조건으로 고르면 디렉터리째 건너뛴다(PartitionFilters). 셋 다 옵티마이저가 조건을 이해할 수 있을 때만 일어난다. 칼럼을 함수로 감싸는 순간 밀어내기는 사라진다.
AQE 는 셔플이 끝난 뒤 실제 크기를 보고 남은 계획을 고친다. 그래서 돌기 전의 계획(isFinalPlan=false)과 돈 뒤의 계획은 다를 수 있고, 진짜로 무엇이 돌았는지는 돈 뒤의 계획이나 이벤트 로그에서 본다.
단계
- /root/spk/plan/convert.py(앱
spk-plan-convert)로 주문을 스키마로 읽어order_date(날짜) 칼럼을 더하고 /root/spk/plan/orders_pq 에 Parquet 으로 쓰세요. - 도우미 /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()하세요. - /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 에 쓰세요. - /root/spk/plan/nopush.py(앱
spk-plan-nopush)로 같은 질문을upper(status) == 'REFUNDED'로 바꿔 /root/spk/plan/out/nopush 에 쓰고, 그 계획의PushedFilters목록을 /root/spk/plan/out/nopush.json 에 쓰세요. - /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 에 쓰세요. - /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": 정수}로 쓰세요. - /root/spk/plan/optimize.py(앱
spk-plan-opt)로qty > 1과qty > 3두where를 잇고1 + 2를three로 고르는 질의의extended계획을 /root/spk/plan/out/optimized.txt 에 저장하고collect()하세요. - /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 · EXPLAIN · Parquet Files · Adaptive Query Execution
CSV 를 Parquet 으로
/root/spk/plan/convert.py 를 앱 이름 spk-plan-convert 로 만들어 /data/shop/orders.csv 를 스키마로 읽고 order_date = to_date(order_ts) 칼럼을 더해 /root/spk/plan/orders_pq 에 Parquet 으로 쓰세요(나누지 않고).
열 형식 파일은 칼럼마다 따로 저장되고 행 그룹마다 최솟값·최댓값 통계를 가집니다. 이 두 가지가 뒤의 모든 단계에서 '덜 읽기' 의 재료입니다. 채점기는 행 수와 칼럼 목록, 파일이 정말 Parquet 인지(첫 네 바이트 PAR1)를 봅니다.
explain 의 모드 — 논리 계획에서 물리 계획까지
도우미 /root/spk/plan/planhelp.py 에 explain 출력을 문자열로 받는 explain(df, mode), formatted 계획에서 이름: 값 줄의 값을 꺼내는 field(plan, name), [A(x,1), B(y)] 를 항목 목록으로 나누는 items(text) 를 두고, /root/spk/plan/explain.py 를 앱 이름 spk-plan-explain 으로 만들어 orders_pq 에서 status == 'refunded' 인 행의 order_id,qty 를 고르는 질의의 extended·formatted 계획을 /root/spk/plan/out/extended.txt·/root/spk/plan/out/formatted.txt 에 저장한 뒤 그 질의를 collect() 하세요.
extended 는 네 절을 차례로 찍습니다 — 파싱된 논리 계획, 분석된 논리 계획, 최적화된 논리 계획, 물리 계획. 분석 단계에서 칼럼 이름이 실제 칼럼(#번호)으로 묶이고, 최적화 단계에서 조건이 스캔 가까이 내려갑니다. formatted 는 그 물리 계획을 연산자별로 풀어 씁니다. 채점기는 formatted 파일이 실제로 돈 계획(이벤트 로그)과 같은지 봅니다.
조건 밀어내기와 칼럼 가지치기
/root/spk/plan/pushdown.py 를 앱 이름 spk-plan-push 로 만들어 orders_pq 에서 status == 'refunded' 이고 qty >= 5 인 행의 order_id,qty,channel 을 /root/spk/plan/out/pushdown 에 머리줄 있는 CSV 로 쓰고, 그 질의의 formatted 계획에 있는 PushedFilters 의 항목 목록과 ReadSchema 값을 /root/spk/plan/out/pushdown.json 에 {"pushed_filters": [문자열…], "read_schema": 문자열} 로 쓰세요.
PushedFilters 는 Parquet 읽기에 넘긴 조건이고 ReadSchema 는 실제로 읽는 칼럼입니다. 고르지 않은 customer_id·product_id 가 ReadSchema 에서 빠지고, 조건에만 쓰인 status 는 남는 것을 보세요. 채점기는 여러분이 적은 목록이 실제로 돈 계획의 것과 한 글자씩 같은지 봅니다.
함수로 감싸면 밀어내기가 사라진다
/root/spk/plan/nopush.py 를 앱 이름 spk-plan-nopush 로 만들어 3단계와 같은 질문을 F.upper("status") == "REFUNDED" 와 qty >= 5 로 바꿔 /root/spk/plan/out/nopush 에 머리줄 있는 CSV 로 쓰고, 그 계획의 PushedFilters 항목 목록을 /root/spk/plan/out/nopush.json 에 {"pushed_filters": [문자열…]} 로 쓰세요.
결과 행은 3단계와 똑같습니다. 달라지는 것은 계획입니다. Parquet 은 upper(status) 가 무엇인지 모르므로 그 조건은 읽기에 넘겨지지 않고, 파일을 다 읽은 뒤 Spark 가 거릅니다. 실무에서 날짜 칼럼을 date_format 으로 감싸 비교하다 이 함정에 자주 빠집니다.
파티션 가지치기 — 디렉터리째 건너뛰기
/root/spk/plan/partition.py 를 앱 이름 spk-plan-part 로 만들어 orders_pq 를 partitionBy("order_date") 로 /root/spk/plan/orders_by_day 에 쓰고, 다시 읽어 order_date == '2026-02-14' 인 행을 센 값과 그 질의 계획의 PartitionFilters 값을 /root/spk/plan/out/partition.json 에 {"rows": 정수, "partition_filters": 문자열} 로 쓰세요.
나눈 칼럼은 파일 안이 아니라 디렉터리 이름(order_date=2026-02-14)에 들어갑니다. 그 칼럼에 거는 조건은 디렉터리를 고르는 데 쓰이고, 고르지 않은 날짜의 파일은 열지도 않습니다. 채점기는 이벤트 로그에서 '읽은 파티션 수' 지표가 1 인지도 봅니다.
AQE 가 돈 뒤에 고친 계획
/root/spk/plan/aqe.py 를 앱 이름 spk-plan-aqe 로 만들어 orders_pq 의 고객별 주문 수(groupBy("customer_id").count())를 collect() 한 뒤에 explain(mode="simple") 출력을 /root/spk/plan/out/aqe_final.txt 에 저장하세요. 그리고 그 앱의 이벤트 로그에서 셔플을 읽은 태스크 수를 세어 /root/spk/plan/out/aqe.json 에 {"shuffle_partitions": 200, "reduce_tasks": 정수} 로 쓰세요.
spark.sql.shuffle.partitions 는 200 이지만 이 자료의 셔플은 몇 MB 뿐이라 AQE 가 파티션을 합칩니다. 돈 뒤의 계획에는 isFinalPlan=true 와 합쳐진 셔플 읽기(AQEShuffleRead)가 보입니다. 셔플을 읽은 태스크 수는 로그의 SparkListenerTaskEnd 에서 Total Records Read 가 0 보다 큰 태스크를 세면 나옵니다.
옵티마이저가 합치고 접는 것
/root/spk/plan/optimize.py 를 앱 이름 spk-plan-opt 로 만들어 orders_pq.where(qty > 1).where(qty > 3).select("order_id", (lit(1) + lit(2)).alias("three")) 의 extended 계획을 /root/spk/plan/out/optimized.txt 에 저장하고 그 질의를 collect() 하세요.
분석된 논리 계획에는 Filter 가 둘이고 (1 + 2) 가 그대로 있습니다. 최적화된 논리 계획에서는 두 조건이 한 Filter 로 합쳐지고(필터 합치기) 1 + 2 는 3 으로 접힙니다(상수 접기). 여러분이 쓴 코드의 모양과 실제로 도는 모양이 다르다는 것을 보는 단계입니다.
덜 읽은 만큼을 숫자로
/root/spk/plan/report.md 에 ## 밀어내기 ## 가지치기 ## AQE 세 절을 쓰세요. 둘째 절에는 5단계의 행 수를, 셋째 절에는 6단계의 셔플 읽기 태스크 수를 숫자로 넣으세요.
첫 절에는 3·4단계에서 PushedFilters 가 어떻게 달라졌는지, 둘째 절에는 파티션 가지치기로 무엇을 건너뛰었는지, 셋째 절에는 200 에서 몇으로 줄었는지를 적으세요.