Apache Spark — The answer to a slow job is in the plan and the event log
Measure and fix skew in clicks where one bot makes 40%
한국어 원문으로 표시합니다.
목표
클릭의 40% 를 봇 한 명이 만든 자료를 사용자 표와 조인해, 셔플 뒤 태스크 하나에 레코드가 몰리는 모양을 이벤트 로그로 잰다. AQE 쏠림 조인·솔팅·뜨거운 키 떼어 내기 세 가지로 풀고, 같은 뜨거운 키가 count 집계에서는 왜 문제가 되지 않는지도 확인한다.
왜 중요한가
"잡이 99% 에서 멈춘다" 는 신고의 대부분은 쏠림이다. 셔플은 키의 해시로 파티션을 정하므로 같은 키는 반드시 한 태스크로 간다. 키 하나가 전체의 40% 면 태스크 하나가 40% 를 떠안고, 나머지 태스크가 다 끝난 뒤에도 그 하나만 돈다. 코어를 더 붙여도 소용없다 — 그 태스크는 쪼개지지 않는다. 쏠림은 시간보다 분포로 보는 편이 정확하다. 스테이지 안에서 태스크별 읽은 레코드 수의 최댓값과 중앙값을 견주면, 기계가 빠르든 느리든 같은 숫자가 나온다. 이 실습도 시간이 아니라 그 비율로 판정한다. 처방은 셋이다. AQE 는 셔플 뒤 실제 크기를 보고 큰 파티션을 여러 조각으로 쪼개고 반대편의 짝을 복제한다(설정만 맞으면 코드 수정이 없다). 솔팅은 뜨거운 키에 난수 꼬리표를 붙여 여러 파티션으로 흩고 반대편을 그 수만큼 불린다. 가장 단순한 것은 뜨거운 키를 떼어 내 따로 처리하는 것이다. 그리고 집계 중에는 쏠림이 원래 문제되지 않는 것도 있다 — 맵 쪽에서 먼저 줄이는 부분 집계가 있기 때문이다.
단계
- /root/spk/skew/common.py 에 클릭·사용자 표를 읽는 함수를 두고, /root/spk/skew/hot.py(앱
spk-skew-hot)로 클릭이 많은 사용자 상위 5명(동수이면user_id오름차순)을 /root/spk/skew/out/hot 에 CSV(user_id,clicks)로 쓰세요. - /root/spk/skew/plain.py(앱
spk-skew-plain, 임계값 -1·AQE 끔·셔플 파티션 8)로 클릭과 사용자를 조인해 세그먼트별 클릭 수·ms 합을 /root/spk/skew/out/by_segment 에 CSV(segment,clicks,ms)로 쓰세요. - 2단계 앱의 로그에서, 셔플을 읽은 스테이지 가운데 태스크별 읽은 레코드 최댓값이 가장 큰 스테이지를 골라 그 스테이지 번호·최댓값·중앙값(낮은 쪽)을 /root/spk/skew/out/skew.json 에 쓰세요.
- /root/spk/skew/aqe.py(앱
spk-skew-aqe, 임계값 -1·AQE 켬·셔플 파티션 16·쏠림 임계 256k·권장 크기 64k)로 같은 결과를 /root/spk/skew/out/by_segment_aqe 에 쓰세요. - /root/spk/skew/salt.py(앱
spk-skew-salt, 2단계와 같은 설정)로 봇의 클릭에만 소금 0–7 을 붙이고 사용자 쪽 봇 행을 여덟 벌로 불려user_id·salt로 조인한 결과를 /root/spk/skew/out/by_segment_salt 에 쓰세요. - /root/spk/skew/split.py(앱
spk-skew-split, 2단계와 같은 설정)로 봇의 클릭은 브로드캐스트 조인, 나머지는 보통 조인을 한 뒤unionByName으로 붙인 결과를 /root/spk/skew/out/by_segment_split 에 쓰세요. - /root/spk/skew/agg.py(앱
spk-skew-agg, AQE 끔·셔플 파티션 8)로 사용자별 클릭 수를 /root/spk/skew/out/per_user 에 Parquet 으로 쓰고, 그 앱에서 셔플을 읽은 스테이지의 태스크별 레코드 최댓값·중앙값(낮은 쪽)을 /root/spk/skew/out/agg_skew.json 에 쓰세요. - /root/spk/skew/report.md 에
## 쏠림의 모양## 세 가지 처방## 부분 집계세 절을 쓰세요. 첫 절에 3단계의 두 숫자, 셋째 절에 7단계의 두 숫자를 넣으세요.
참고
- 원본:
/data/clicks/clicks.jsonl(user_id, page, ts, ms — 24만 줄),/data/clicks/users.csv(user_id, segment — 봇은 segmentbot). - 태스크별 읽은 레코드는
SparkListenerTaskEnd의Task Metrics.Shuffle Read Metrics.Total Records Read입니다. 중앙값(낮은 쪽)은 오름차순으로 늘어놓은 n 개 가운데(n-1)//2번째(0부터) 값입니다. - 스크립트는
/root/spk/skew에 두고 거기서 돌리세요. 로그를 읽는 작은 도구(skewtool.py)를 만들어 두면 3·7단계가 한 줄입니다. - 흔한 실수: 2단계에서 AQE 나 브로드캐스트를 켜 둬 쏠림이 사라지는 것, 소금을 양쪽에 난수로 붙여 짝이 안 맞는 것(반대편은 불려야 합니다), 세그먼트 결과가 처방마다 달라지는 것(답은 같아야 합니다).
- 공식 문서: Optimizing Skew Join · Adaptive Query Execution · Monitoring — Executor Task Metrics · Built-in Functions
뜨거운 키 찾기
/root/spk/skew/common.py 에 클릭(/data/clicks/clicks.jsonl)과 사용자(/data/clicks/users.csv) 표를 읽는 함수를 두고, /root/spk/skew/hot.py 를 앱 이름 spk-skew-hot 으로 만들어 클릭 수 상위 5명(클릭 수 내림차순, 같으면 user_id 오름차순)을 /root/spk/skew/out/hot 에 머리줄 있는 CSV(user_id,clicks)로 쓰세요.
쏠림을 푸는 첫걸음은 어느 키가 뜨거운지 아는 것입니다. 1등과 2등의 차이를 보세요. 1등이 나머지와 몇 배 차이가 나면, 그 키 하나가 한 태스크를 붙잡습니다.
그냥 조인하면 태스크 하나에 몰린다
/root/spk/skew/plain.py 를 앱 이름 spk-skew-plain, 설정 spark.sql.autoBroadcastJoinThreshold=-1·spark.sql.adaptive.enabled=false·spark.sql.shuffle.partitions=8 로 만들어 클릭과 사용자를 user_id 로 조인하고 세그먼트별 클릭 수(clicks)와 ms 합(ms)을 /root/spk/skew/out/by_segment 에 머리줄 있는 CSV(segment,clicks,ms)로 쓰세요.
브로드캐스트와 AQE 를 끈 것은 쏠림을 일부러 드러내기 위해서입니다. 정렬 병합 조인은 양쪽을 user_id 로 셔플하므로 봇의 클릭 9만여 줄이 전부 한 파티션으로 갑니다. 채점기는 그 스테이지에서 최댓값이 중앙값의 몇 배인지를 봅니다.
쏠림을 숫자로 — 최댓값과 중앙값
2단계 앱(spk-skew-plain)의 가장 최근 로그에서, 셔플을 읽은 스테이지마다 태스크별 Total Records Read 를 모아 최댓값이 가장 큰 스테이지를 고르고 /root/spk/skew/out/skew.json 에 {"stage_id": 정수, "max_records": 정수, "median_records": 정수} 로 쓰세요. 중앙값은 오름차순 n 개 중 (n-1)//2 번째(0부터) 값입니다.
조인 스테이지는 양쪽 셔플을 함께 읽으므로 봇이 든 파티션의 태스크는 봇 클릭 전부와 봇 사용자 한 줄을 읽습니다. 최댓값 ÷ 중앙값이 쏠림의 크기입니다. 이 비율이 3 을 넘으면 흔히 '쏠렸다' 고 부릅니다.
처방 1 — AQE 쏠림 조인
/root/spk/skew/aqe.py 를 앱 이름 spk-skew-aqe, 설정 spark.sql.autoBroadcastJoinThreshold=-1·spark.sql.shuffle.partitions=16·spark.sql.adaptive.skewJoin.skewedPartitionThresholdInBytes=256k·spark.sql.adaptive.advisoryPartitionSizeInBytes=64k(AQE 는 켠 채)로 만들어 2단계와 같은 결과를 /root/spk/skew/out/by_segment_aqe 에 쓰세요.
AQE 는 파티션이 중앙값의 몇 배(skewedPartitionFactor, 기본 5)이면서 임계 바이트보다 클 때 쏠렸다고 봅니다. 기본 임계(256MB)는 이 작은 자료에 맞지 않아 낮췄습니다. 최종 계획에 SortMergeJoin(skew=true) 가 보이면 코드 수정 없이 쪼갠 것입니다.
처방 2 — 뜨거운 키에 소금 뿌리기
/root/spk/skew/salt.py 를 앱 이름 spk-skew-salt, 2단계와 같은 설정(임계값 -1·AQE 끔·셔플 파티션 8)으로 만들어, 클릭 쪽은 봇이면 0–7 사이 난수 salt, 아니면 0 을 붙이고, 사용자 쪽은 봇 행만 salt 0–7 여덟 벌로 불린 뒤(나머지는 0) user_id·salt 로 조인한 결과를 /root/spk/skew/out/by_segment_salt 에 쓰세요.
조인 키에 소금이 들어가니 봇의 클릭이 여덟 파티션으로 흩어집니다. 반대편 봇 행이 여덟 벌이어야 어느 소금으로 가든 짝이 있습니다. 결과는 2단계와 똑같아야 하고, 채점기는 쏠림 비율이 2단계의 절반 아래로 떨어졌는지 봅니다.
처방 3 — 뜨거운 키 떼어 내기
/root/spk/skew/split.py 를 앱 이름 spk-skew-split, 2단계와 같은 설정으로 만들어 봇의 클릭은 봇 사용자 한 줄과 broadcast 조인하고, 나머지 클릭은 보통 조인한 뒤 unionByName 으로 붙여 /root/spk/skew/out/by_segment_split 에 쓰세요.
뜨거운 키가 하나이고 이름을 안다면 가장 단순하고 확실한 처방입니다. 그 키의 반대편은 한 줄이라 브로드캐스트가 공짜이고, 나머지는 고르게 퍼집니다. 계획에 BroadcastHashJoin·SortMergeJoin·Union 이 함께 보여야 합니다.
count 는 왜 쏠리지 않나 — 부분 집계
/root/spk/skew/agg.py 를 앱 이름 spk-skew-agg, 설정 spark.sql.adaptive.enabled=false·spark.sql.shuffle.partitions=8 로 만들어 사용자별 클릭 수(user_id,clicks)를 /root/spk/skew/out/per_user 에 Parquet 으로 쓰고, 그 앱에서 셔플을 읽은 스테이지의 태스크별 레코드 최댓값·중앙값(낮은 쪽)을 /root/spk/skew/out/agg_skew.json 에 {"max_records": 정수, "median_records": 정수} 로 쓰세요.
같은 봇이 있는데 이번에는 태스크가 고르게 나뉩니다. 계획의 첫 HashAggregate 는 partial_count — 맵 태스크마다 사용자별로 먼저 세어 두므로, 셔플로 넘어가는 것은 맵 태스크당 사용자 한 줄입니다. 쏠림이 문제되는 것은 이렇게 먼저 줄일 수 없는 연산(조인, collect_list 같은 것)입니다.
쏠림과 처방을 숫자로 남기기
/root/spk/skew/report.md 에 ## 쏠림의 모양 ## 세 가지 처방 ## 부분 집계 세 절을 쓰세요. 첫 절에 3단계의 최댓값·중앙값을, 셋째 절에 7단계의 최댓값·중앙값을 숫자로 넣으세요.
둘째 절에는 세 처방이 각각 무엇을 바꾸는지(코드인가 설정인가, 무엇을 복제하는가)를 한 줄씩 적으세요. 운영에서 어느 것을 먼저 시도할지도 적으면 좋습니다.