Apache Spark — The answer to a slow job is in the plan and the event log
Use explain to confirm how much less Parquet is read
한국어 원문으로 표시합니다.
목표
주문을 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 에서 몇으로 줄었는지를 적으세요.