Apache Spark — 느린 잡의 답은 실행 계획과 이벤트 로그에 있다 · 쏠림 · 实验
봇 하나가 40% 인 클릭으로 쏠림을 재고 푼다
목표
클릭의 40% 를 봇 한 명이 만든 자료를 사용자 표와 조인해, 셔플 뒤 태스크 하나에 레코드가 몰리는 모양을 이벤트 로그로 잰다. AQE 쏠림 조인·솔팅·뜨거운 키 떼어 내기 세 가지로 풀고, 같은 뜨거운 키가 count 집계에서는 왜 문제가 되지 않는지도 확인한다.
왜 중요한가
"잡이 99% 에서 멈춘다" 는 신고의 대부분은 쏠림이다. 셔플은 키의 해시로 파티션을 정하므로 같은 키는 반드시 한 태스크로 간다. 키 하나가 전체의 40% 면 태스크 하나가 40% 를 떠안고, 나머지 태스크가 다 끝난 뒤에도 그 하나만 돈다. 코어를 더 붙여도 소용없다 — 그 태스크는 쪼개지지 않는다.
쏠림은 시간보다 분포로 보는 편이 정확하다. 스테이지 안에서 태스크별 읽은 레코드 수의 최댓값과 중앙값을 견주면, 기계가 빠르든 느리든 같은 숫자가 나온다. 이 실습도 시간이 아니라 그 비율로 판정한다.
처방은 셋이다. AQE 는 셔플 뒤 실제 크기를 보고 큰 파티션을 여러 조각으로 쪼개고 반대편의 짝을 복제한다(설정만 맞으면 코드 수정이 없다). 솔팅은 뜨거운 키에 난수 꼬리표를 붙여 여러 파티션으로 흩고 반대편을 그 수만큼 불린다. 가장 단순한 것은 뜨거운 키를 떼어 내 따로 처리하는 것이다. 그리고 집계 중에는 쏠림이 원래 문제되지 않는 것도 있다 — 맵 쪽에서 먼저 줄이는 부분 집계가 있기 때문이다.
단계
1. /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)로 쓰세요.
2. /root/spk/skew/plain.py(앱 spk-skew-plain, 임계값 -1·AQE 끔·셔플 파티션 8)로 클릭과 사용자를 조인해 세그먼트별 클릭 수·ms 합을 /root/spk/skew/out/by_segment 에 CSV(segment,clicks,ms)로 쓰세요.
3. 2단계 앱의 로그에서, 셔플을 읽은 스테이지 가운데 태스크별 읽은 레코드 최댓값이 가장 큰 스테이지를 골라 그 스테이지 번호·최댓값·중앙값(낮은 쪽)을 /root/spk/skew/out/skew.json 에 쓰세요.
4. /root/spk/skew/aqe.py(앱 spk-skew-aqe, 임계값 -1·AQE 켬·셔플 파티션 16·쏠림 임계 256k·권장 크기 64k)로 같은 결과를 /root/spk/skew/out/by_segment_aqe 에 쓰세요.
5. /root/spk/skew/salt.py(앱 spk-skew-salt, 2단계와 같은 설정)로 봇의 클릭에만 소금 0–7 을 붙이고 사용자 쪽 봇 행을 여덟 벌로 불려 user_id·salt 로 조인한 결과를 /root/spk/skew/out/by_segment_salt 에 쓰세요.
6. /root/spk/skew/split.py(앱 spk-skew-split, 2단계와 같은 설정)로 봇의 클릭은 브로드캐스트 조인, 나머지는 보통 조인을 한 뒤 unionByName 으로 붙인 결과를 /root/spk/skew/out/by_segment_split 에 쓰세요.
7. /root/spk/skew/agg.py(앱 spk-skew-agg, AQE 끔·셔플 파티션 8)로 사용자별 클릭 수를 /root/spk/skew/out/per_user 에 Parquet 으로 쓰고, 그 앱에서 셔플을 읽은 스테이지의 태스크별 레코드 최댓값·중앙값(낮은 쪽)을 /root/spk/skew/out/agg_skew.json 에 쓰세요.
8. /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](https://spark.apache.org/docs/4.2.0/sql-performance-tuning.html#optimizing-skew-join) · [Adaptive Query Execution](https://spark.apache.org/docs/4.2.0/sql-performance-tuning.html#adaptive-query-execution) · [Monitoring — Executor Task Metrics](https://spark.apache.org/docs/4.2.0/monitoring.html#executor-task-metrics) · [Built-in Functions](https://spark.apache.org/docs/4.2.0/sql-ref-functions-builtin.html)
8个步骤
- 뜨거운 키 찾기
- 그냥 조인하면 태스크 하나에 몰린다
- 쏠림을 숫자로 — 최댓값과 중앙값
- 처방 1 — AQE 쏠림 조인
- 처방 2 — 뜨거운 키에 소금 뿌리기
- 처방 3 — 뜨거운 키 떼어 내기
- count 는 왜 쏠리지 않나 — 부분 집계
- 쏠림과 처방을 숫자로 남기기