LabHub
배우기 러닝패스 코스

Apache Spark — 遅いジョブの答えは実行計画とイベントログにある

explain で Parquet をどれだけ読まずに済むか確かめる

LabHub 에서 이어서 보기

한국어 원문으로 표시합니다.

목표

주문을 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.pyexplain·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_dayorder_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 > 1qty > 3where 를 잇고 1 + 2three 로 고르는 질의의 extended 계획을 /root/spk/plan/out/optimized.txt 에 저장하고 collect() 하세요.
  8. /root/spk/plan/report.md## 밀어내기 ## 가지치기 ## AQE 세 절을 쓰세요. 둘째 절에는 5단계의 행 수를, 셋째 절에는 6단계의 태스크 수를 숫자로 넣으세요.

참고

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_pqpartitionBy("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 + 23 으로 접힙니다(상수 접기). 여러분이 쓴 코드의 모양과 실제로 도는 모양이 다르다는 것을 보는 단계입니다.

덜 읽은 만큼을 숫자로

/root/spk/plan/report.md## 밀어내기 ## 가지치기 ## AQE 세 절을 쓰세요. 둘째 절에는 5단계의 행 수를, 셋째 절에는 6단계의 셔플 읽기 태스크 수를 숫자로 넣으세요.

첫 절에는 3·4단계에서 PushedFilters 가 어떻게 달라졌는지, 둘째 절에는 파티션 가지치기로 무엇을 건너뛰었는지, 셋째 절에는 200 에서 몇으로 줄었는지를 적으세요.