LabHub
배우기 러닝패스 코스

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

LabHub 에서 이어서 보기

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

목표

주문과 상품의 조인을 브로드캐스트 해시·정렬 병합·셔플 해시 세 전략으로 돌려 계획과 결과를 견주고, 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단계의 두 숫자를 넣으세요.

참고

작은 쪽은 알아서 브로드캐스트된다

/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단계의 두 숫자를 넣으세요.

첫 절에는 각 전략이 무엇을 셔플하고 무엇을 메모리에 올리는지를, 둘째 절에는 처음 계획과 최종 계획이 어떻게 달랐는지를, 셋째 절에는 행이 몇 개 불었는지를 적으세요.