Apache Spark — The answer to a slow job is in the plan and the event log
Run the same join with four strategies and confirm each in the plan
한국어 원문으로 표시합니다.
목표
주문과 상품의 조인을 브로드캐스트 해시·정렬 병합·셔플 해시 세 전략으로 돌려 계획과 결과를 견주고, AQE 가 실행 중에 정렬 병합을 브로드캐스트로 바꾸는 것을 이벤트 로그로 확인한다. 키가 겹치는 표와의 조인이 행을 불리는 것과, 반대편에 짝이 없는 행을 고르는 조인도 해 본다.
왜 중요한가
조인은 Spark 잡에서 가장 비싼 연산인 경우가 많고, 그 비용은 전략이 정한다. 한쪽이 작으면 작은 쪽을 모든 태스크에 복사해(브로드캐스트) 큰 쪽을 셔플하지 않고 조인한다. 둘 다 크면 양쪽을 키로 셔플한 뒤 정렬해 맞춰 가며 합친다(정렬 병합). 셔플 뒤 한쪽을 해시 표로 만들어 조인하는 셔플 해시는 정렬을 건너뛰는 대신 메모리를 쓴다.
Spark 는 통계로 크기를 짐작해 전략을 고른다. 작다고 보는 기준이 spark.sql.autoBroadcastJoinThreshold(기본 10MB)다. 짐작이 틀리면 전략도 틀린다 — 필터 뒤 실제로는 작은데 원본 크기로 짐작해 정렬 병합을 고르거나, 반대로 큰 표를 브로드캐스트하다 드라이버 메모리가 넘친다. AQE 는 셔플이 끝난 뒤의 실제 크기로 다시 판단해 이것을 고친다.
조인 결과의 행 수는 키의 중복이 정한다. 한쪽 키가 유일하다고 믿었는데 아니면, 조인은 오류 없이 행을 곱해 버리고 뒤의 합계가 전부 부푼다.
단계
- /root/spk/join/common.py 에 네 표를 읽는 함수와 분류별 매출 함수를 두고, /root/spk/join/auto.py(앱
spk-join-auto, 기본 설정)로 분류별 매출을 /root/spk/join/out/by_category 에 머리줄 있는 CSV(category,revenue)로 쓰세요. - /root/spk/join/smj.py(앱
spk-join-smj,spark.sql.autoBroadcastJoinThreshold=-1·spark.sql.adaptive.enabled=false)로 같은 결과를 /root/spk/join/out/by_category_smj 에 쓰세요. - /root/spk/join/hint.py(앱
spk-join-hint, 2단계와 같은 설정)에서 상품 쪽에broadcast힌트를 주어 /root/spk/join/out/by_category_hint 에 쓰세요. - /root/spk/join/shash.py(앱
spk-join-shash, 2단계와 같은 설정)에서 상품 쪽에shuffle_hash힌트를 주어 /root/spk/join/out/by_category_shash 에 쓰세요. - /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)로 쓰세요. - /root/spk/join/dup.py(앱
spk-join-dup)로 결제 완료 주문을 판촉 표와 그냥 조인한 행 수와left_semi로 조인한 행 수를 /root/spk/join/out/dup.json 에{"naive": 정수, "semi": 정수}로 쓰세요. - /root/spk/join/anti.py(앱
spk-join-anti)로 한 번도 주문하지 않은 고객의customer_id를 /root/spk/join/out/no_orders 에 CSV 로 쓰세요. - /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 · Hints · JOIN · Converting sort-merge join to broadcast join
작은 쪽은 알아서 브로드캐스트된다
/root/spk/join/common.py 에 네 표(주문·상품·고객·판촉)를 스키마로 읽는 함수와 분류별 매출(결제 완료 주문 × 상품, category,revenue=qty×price 합) 함수를 두고, /root/spk/join/auto.py 를 앱 이름 spk-join-auto(설정 기본값)로 만들어 결과를 /root/spk/join/out/by_category 에 머리줄 있는 CSV 로 쓰세요.
상품 표는 수 KB 라 임계값(10MB)보다 훨씬 작습니다. Spark 는 이것을 모든 태스크에 복사하고 주문 쪽은 셔플하지 않습니다. 채점기는 이벤트 로그의 계획에 BroadcastHashJoin 이 있는지, 분류별 매출이 원본과 같은지를 봅니다.
임계값을 끄면 정렬 병합
/root/spk/join/smj.py 를 앱 이름 spk-join-smj, 설정 spark.sql.autoBroadcastJoinThreshold=-1·spark.sql.adaptive.enabled=false 로 만들어 1단계와 같은 결과를 /root/spk/join/out/by_category_smj 에 쓰세요.
임계값 -1 은 '자동 브로드캐스트 없음' 입니다. 이제 주문과 상품 양쪽이 product_id 로 셔플되고 정렬된 뒤 합쳐집니다. 결과는 1단계와 한 줄도 다르지 않아야 합니다. AQE 를 끄는 이유는 실행 중에 전략이 바뀌지 않게 하려는 것입니다.
힌트로 브로드캐스트를 강제하기
/root/spk/join/hint.py 를 앱 이름 spk-join-hint, 2단계와 같은 설정(임계값 -1, AQE 끔)으로 만들되 상품 쪽을 F.broadcast(p) 로 감싸 /root/spk/join/out/by_category_hint 에 쓰세요.
힌트는 통계보다 앞섭니다. 임계값을 꺼 두었어도 힌트가 있으면 브로드캐스트합니다. 반대로 말하면, 큰 표에 무심코 붙인 브로드캐스트 힌트는 드라이버와 모든 실행기의 메모리를 그대로 먹습니다.
셔플 해시 조인 — 정렬을 건너뛰는 대신
/root/spk/join/shash.py 를 앱 이름 spk-join-shash, 2단계와 같은 설정으로 만들되 상품 쪽에 p.hint("shuffle_hash") 를 주어 /root/spk/join/out/by_category_shash 에 쓰세요.
셔플 해시 조인은 양쪽을 키로 셔플한 뒤, 작은 쪽 파티션을 해시 표로 만들어 큰 쪽을 흘려 보냅니다. 정렬이 없어 빠를 수 있지만 해시 표가 메모리에 들어가야 합니다. 계획에서 Sort 연산자가 사라진 것을 확인하세요.
AQE 가 실행 중에 전략을 바꾼다
/root/spk/join/aqe.py 를 앱 이름 spk-join-aqe, 설정 spark.sql.autoBroadcastJoinThreshold=100k(AQE 는 켜 둔 채)로 만들어 결제 완료 주문을 tier == 'vip' 인 고객과 customer_id 로 조인하고 도시별 주문 수를 /root/spk/join/out/vip_by_city 에 머리줄 있는 CSV(city,orders)로 쓰세요.
고객 표 파일은 100KB 보다 크니 처음 계획은 정렬 병합입니다. 하지만 vip 로 거른 뒤의 실제 크기는 수십 KB 입니다. AQE 는 셔플 맵 단계가 끝난 뒤 그 크기를 보고 남은 계획을 브로드캐스트로 바꿉니다. 채점기는 같은 실행의 처음 계획과 최종 계획을 견줍니다.
키가 겹치면 조인이 행을 불린다
/root/spk/join/dup.py 를 앱 이름 spk-join-dup 으로 만들어 결제 완료 주문을 판촉 표(/data/shop/promos.csv)와 product_id 로 그냥 조인한 행 수와 left_semi 로 조인한 행 수를 /root/spk/join/out/dup.json 에 {"naive": 정수, "semi": 정수} 로 쓰세요.
판촉 표에는 코드가 둘인 상품이 있습니다. 그 상품의 주문은 그냥 조인하면 두 줄이 됩니다. left_semi 는 '짝이 있는가' 만 보므로 왼쪽 행을 늘리지 않습니다. 매출을 판촉 여부로 나눌 때 어느 쪽을 써야 할지 생각해 보세요.
짝이 없는 쪽 고르기 — left_anti
/root/spk/join/anti.py 를 앱 이름 spk-join-anti 로 만들어 한 번도 주문하지 않은(상태와 무관) 고객의 customer_id 를 /root/spk/join/out/no_orders 에 머리줄 있는 CSV 로 쓰세요.
not in 서브쿼리나 left join 뒤 null 거르기로도 되지만, left_anti 는 뜻이 그대로 보이고 null 이 섞인 키에서도 헷갈리지 않습니다. 계획에서 조인 종류가 LeftAnti 로 찍히는지 보세요.
어느 전략이 언제 맞는지 남기기
/root/spk/join/report.md 에 ## 네 가지 전략 ## AQE 의 전환 ## 키 중복 세 절을 쓰세요. 첫 절에는 1–4단계에서 본 조인 연산자 이름을, 셋째 절에는 6단계의 두 숫자를 넣으세요.
첫 절에는 각 전략이 무엇을 셔플하고 무엇을 메모리에 올리는지를, 둘째 절에는 처음 계획과 최종 계획이 어떻게 달랐는지를, 셋째 절에는 행이 몇 개 불었는지를 적으세요.