Apache Spark — 느린 잡의 답은 실행 계획과 이벤트 로그에 있다 · 조인 전략 · 实验
같은 조인을 네 가지 전략으로 돌리고 계획으로 확인한다
목표
주문과 상품의 조인을 브로드캐스트 해시·정렬 병합·셔플 해시 세 전략으로 돌려 계획과 결과를 견주고, AQE 가 실행 중에 정렬 병합을 브로드캐스트로 바꾸는 것을 이벤트 로그로 확인한다. 키가 겹치는 표와의 조인이 행을 불리는 것과, 반대편에 짝이 없는 행을 고르는 조인도 해 본다.
왜 중요한가
조인은 Spark 잡에서 가장 비싼 연산인 경우가 많고, 그 비용은 전략이 정한다. 한쪽이 작으면 작은 쪽을 모든 태스크에 복사해(브로드캐스트) 큰 쪽을 셔플하지 않고 조인한다. 둘 다 크면 양쪽을 키로 셔플한 뒤 정렬해 맞춰 가며 합친다(정렬 병합). 셔플 뒤 한쪽을 해시 표로 만들어 조인하는 셔플 해시는 정렬을 건너뛰는 대신 메모리를 쓴다.
Spark 는 통계로 크기를 짐작해 전략을 고른다. 작다고 보는 기준이 spark.sql.autoBroadcastJoinThreshold(기본 10MB)다. 짐작이 틀리면 전략도 틀린다 — 필터 뒤 실제로는 작은데 원본 크기로 짐작해 정렬 병합을 고르거나, 반대로 큰 표를 브로드캐스트하다 드라이버 메모리가 넘친다. AQE 는 셔플이 끝난 뒤의 실제 크기로 다시 판단해 이것을 고친다.
조인 결과의 행 수는 키의 중복이 정한다. 한쪽 키가 유일하다고 믿었는데 아니면, 조인은 오류 없이 행을 곱해 버리고 뒤의 합계가 전부 부푼다.
단계
1. /root/spk/join/common.py 에 네 표를 읽는 함수와 분류별 매출 함수를 두고, /root/spk/join/auto.py(앱 spk-join-auto, 기본 설정)로 분류별 매출을 /root/spk/join/out/by_category 에 머리줄 있는 CSV(category,revenue)로 쓰세요.
2. /root/spk/join/smj.py(앱 spk-join-smj, spark.sql.autoBroadcastJoinThreshold=-1·spark.sql.adaptive.enabled=false)로 같은 결과를 /root/spk/join/out/by_category_smj 에 쓰세요.
3. /root/spk/join/hint.py(앱 spk-join-hint, 2단계와 같은 설정)에서 상품 쪽에 broadcast 힌트를 주어 /root/spk/join/out/by_category_hint 에 쓰세요.
4. /root/spk/join/shash.py(앱 spk-join-shash, 2단계와 같은 설정)에서 상품 쪽에 shuffle_hash 힌트를 주어 /root/spk/join/out/by_category_shash 에 쓰세요.
5. /root/spk/join/aqe.py(앱 spk-join-aqe, spark.sql.autoBroadcastJoinThreshold=100k, AQE 켬)로 결제 완료 주문을 tier == 'vip' 고객과 조인해 도시별 주문 수를 /root/spk/join/out/vip_by_city 에 CSV(city,orders)로 쓰세요.
6. /root/spk/join/dup.py(앱 spk-join-dup)로 결제 완료 주문을 판촉 표와 그냥 조인한 행 수와 left_semi 로 조인한 행 수를 /root/spk/join/out/dup.json 에 {"naive": 정수, "semi": 정수} 로 쓰세요.
7. /root/spk/join/anti.py(앱 spk-join-anti)로 한 번도 주문하지 않은 고객의 customer_id 를 /root/spk/join/out/no_orders 에 CSV 로 쓰세요.
8. /root/spk/join/report.md 에 ## 네 가지 전략 ## AQE 의 전환 ## 키 중복 세 절을 쓰세요. 셋째 절에는 6단계의 두 숫자를 넣으세요.
참고
- 원본: 주문
/data/shop/orders.csv, 상품/data/shop/products.csv(400행), 고객/data/shop/customers.csv(2만 행), 판촉/data/shop/promos.csv(product_id, promo_code — 상품 몇 개는 코드가 둘). - 전략은 계획의 연산자 이름으로 확인합니다:
BroadcastHashJoin,SortMergeJoin,ShuffledHashJoin,BroadcastHashJoin … LeftSemi,… LeftAnti. AQE 가 돈 뒤의 계획은== Final Plan ==절에 있습니다. - 힌트는
F.broadcast(df)나df.hint("broadcast")·df.hint("shuffle_hash")·df.hint("merge")로 줍니다. SQL 이면/*+ BROADCAST(p) */. - 흔한 실수: 2–4단계에서 AQE 를 켜 두어 계획이 실행 중에 바뀌는 것, 판촉 표처럼 키가 겹치는 표를 차원 표로 믿는 것.
- 공식 문서: [Join Strategy Hints](https://spark.apache.org/docs/4.2.0/sql-performance-tuning.html#join-strategy-hints-for-sql-queries) · [Hints](https://spark.apache.org/docs/4.2.0/sql-ref-syntax-qry-select-hints.html) · [JOIN](https://spark.apache.org/docs/4.2.0/sql-ref-syntax-qry-select-join.html) · [Converting sort-merge join to broadcast join](https://spark.apache.org/docs/4.2.0/sql-performance-tuning.html#converting-sort-merge-join-to-broadcast-join)
8个步骤
- 작은 쪽은 알아서 브로드캐스트된다
- 임계값을 끄면 정렬 병합
- 힌트로 브로드캐스트를 강제하기
- 셔플 해시 조인 — 정렬을 건너뛰는 대신
- AQE 가 실행 중에 전략을 바꾼다
- 키가 겹치면 조인이 행을 불린다
- 짝이 없는 쪽 고르기 — left_anti
- 어느 전략이 언제 맞는지 남기기