LabHub
배우기 러닝패스 코스

Apache Spark — 遅いジョブの答えは実行計画とイベントログにある

シャッフルパーティション 200 が生んだファイルとタスクを減らす

LabHub 에서 이어서 보기

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

목표

고객별 매출 집계를 AQE 를 끈 기본값(셔플 파티션 200)·손으로 줄인 값(8)·AQE 로 세 번 돌려, 셔플 뒤 태스크 수와 결과 파일 수가 어떻게 달라지는지 이벤트 로그로 잰다. repartition 과 coalesce 가 셔플 여부와 파일 수에서 어떻게 다른지도 확인한다.

왜 중요한가

groupBy·join 같은 넓은 변환은 같은 키를 한 태스크로 모아야 하므로 셔플을 만든다. 앞 스테이지가 결과를 파티션별 파일로 쓰고, 뒤 스테이지가 그 조각들을 모아 읽는다. 뒤 스테이지의 태스크 수는 spark.sql.shuffle.partitions 가 정하고 기본은 200 이다. 200 은 큰 클러스터를 기준으로 한 숫자다. 몇 MB 자료에 200 을 쓰면 태스크 200개가 각자 조금씩 일하고, 결과를 파일로 쓰면 작은 파일이 200개 생긴다. 반대로 수백 GB 에 200 이면 태스크 하나가 수 GB 를 떠안고 디스크로 흘린다. 그래서 이 숫자는 자료 크기에 맞춰야 하고, AQE 는 셔플이 끝난 뒤 실제 크기를 보고 작은 조각들을 합쳐 이 일을 대신 한다. 파티션 수를 바꾸는 방법은 둘이다. repartition(n) 은 셔플로 조각을 새로 만들고, coalesce(n) 은 셔플 없이 이웃한 조각을 붙인다. coalesce 는 싸지만 줄이기만 할 수 있다 — 조각 둘을 넷으로 만들 수는 없다.

단계

  1. /root/spk/shuffle/common.py 에 고객별 매출 DataFrame 을 만드는 함수를 두고, /root/spk/shuffle/noaqe.py(앱 spk-shuffle-noaqe, spark.sql.adaptive.enabled=false)로 결과를 /root/spk/shuffle/out/by_customer_200 에 Parquet 으로 쓰세요.
  2. 그 결과 폴더의 데이터 파일(part- 로 시작) 수를 세어 /root/spk/shuffle/out/files_200.txt 에 정수로 적으세요.
  3. /root/spk/shuffle/eight.py(앱 spk-shuffle-8, AQE 끔, spark.sql.shuffle.partitions=8)로 같은 결과를 /root/spk/shuffle/out/by_customer_8 에 쓰세요.
  4. /root/spk/shuffle/aqe.py(앱 spk-shuffle-aqe, 설정 기본값)로 같은 결과를 /root/spk/shuffle/out/by_customer_aqe 에 쓰고, 그 앱에서 셔플을 읽은 태스크 수를 /root/spk/shuffle/out/aqe_tasks.txt 에 정수로 적으세요.
  5. 클릭 원본을 /root/spk/shuffle/repart.py(앱 spk-shuffle-repart)로 repartition(4) 해서 /root/spk/shuffle/out/clicks_repart 에, /root/spk/shuffle/coalesce.py(앱 spk-shuffle-coalesce)로 coalesce(4) 해서 /root/spk/shuffle/out/clicks_coalesce 에 Parquet 으로 쓰고, 두 폴더의 데이터 파일 수를 /root/spk/shuffle/out/repart.json{"repartition_files": 정수, "coalesce_files": 정수} 로 쓰세요.
  6. 3단계 앱이 셔플로 쓴 바이트 합을 로그에서 구해 /root/spk/shuffle/out/shuffle_bytes.json{"app": "spk-shuffle-8", "shuffle_write_bytes": 정수} 로 쓰세요.
  7. /root/spk/shuffle/narrow.py(앱 spk-shuffle-narrow)로 결제 완료 주문에 day 칼럼을 더해 order_id,customer_id,qty,day/root/spk/shuffle/out/narrow 에 쓰세요. 이 앱에는 셔플이 없어야 합니다.
  8. /root/spk/shuffle/report.md## 200 개의 파티션 ## AQE 가 합친 것 ## repartition 과 coalesce 세 절을 쓰세요. 각 절에 2·4·5단계의 숫자를 넣으세요.

참고

AQE 를 끄고 기본값 200 으로

/root/spk/shuffle/common.py 에 결제 완료 주문의 고객별 매출(customer_id, revenue=qty×price 합) DataFrame 을 만드는 함수를 두고, /root/spk/shuffle/noaqe.py 를 앱 이름 spk-shuffle-noaqe, 설정 spark.sql.adaptive.enabled=false 로 만들어 결과를 /root/spk/shuffle/out/by_customer_200 에 Parquet 으로 쓰세요.

AQE 를 끄면 셔플 뒤 스테이지는 정확히 spark.sql.shuffle.partitions 개의 태스크로 돕니다. 채점기는 로그에서 그 스테이지의 태스크 수가 200 인지, 결과가 원본에서 계산한 고객별 매출과 같은지를 봅니다.

200 이 만든 파일 세기

/root/spk/shuffle/out/by_customer_200 안의 데이터 파일(part- 로 시작하는 것) 수를 세어 /root/spk/shuffle/out/files_200.txt 에 정수 하나로 적으세요.

셔플 뒤 태스크 하나가 파일 하나를 씁니다(빈 파티션은 파일을 쓰지 않을 수 있습니다). 고객 2만 명의 매출은 수백 KB 인데 파일이 몇 개인지 보세요 — 이것이 작은 파일 문제의 가장 흔한 출처입니다. _SUCCESS.crc 는 세지 않습니다.

손으로 8 로 줄이기

/root/spk/shuffle/eight.py 를 앱 이름 spk-shuffle-8, 설정 spark.sql.adaptive.enabled=false·spark.sql.shuffle.partitions=8 로 만들어 같은 결과를 /root/spk/shuffle/out/by_customer_8 에 Parquet 으로 쓰세요.

이번에는 셔플 뒤 태스크가 8개, 파일도 8개 이하입니다. 결과는 1단계와 한 줄도 다르지 않아야 합니다 — 파티션 수는 일을 나누는 방법이지 답이 아닙니다.

AQE 가 셔플 뒤에 합친다

/root/spk/shuffle/aqe.py 를 앱 이름 spk-shuffle-aqe 로(설정은 기본값 그대로) 만들어 같은 결과를 /root/spk/shuffle/out/by_customer_aqe 에 쓰고, 그 앱에서 셔플을 읽은 태스크 수를 로그에서 세어 /root/spk/shuffle/out/aqe_tasks.txt 에 정수로 적으세요.

AQE 는 셔플 맵 쪽이 끝난 뒤 파티션별 실제 크기를 보고, 목표 크기(spark.sql.adaptive.advisoryPartitionSizeInBytes)에 못 미치는 이웃 조각을 합칩니다. 셔플 파티션은 여전히 200 인데 읽는 태스크는 훨씬 적어집니다. 최종 계획에는 AQEShuffleRead … coalesced 가 보입니다.

repartition 과 coalesce — 셔플이 있고 없고

/root/spk/shuffle/repart.py(앱 spk-shuffle-repart)로 /data/clicks/clicks.jsonlrepartition(4) 해서 /root/spk/shuffle/out/clicks_repart 에, /root/spk/shuffle/coalesce.py(앱 spk-shuffle-coalesce)로 같은 원본을 coalesce(4) 해서 /root/spk/shuffle/out/clicks_coalesce 에 Parquet 으로 쓰세요. 두 폴더의 데이터 파일 수를 /root/spk/shuffle/out/repart.json{"repartition_files": 정수, "coalesce_files": 정수} 로 쓰세요.

18MB 짜리 원본은 local[2] 에서 몇 조각으로 읽히는지 먼저 보세요(rdd.getNumPartitions()). repartition 은 셔플로 정확히 넷을 만들고, coalesce 는 셔플 없이 이웃을 붙일 뿐이라 원래 조각 수보다 늘지 않습니다. 채점기는 repart 앱에만 셔플 쓰기가 있는지도 로그로 봅니다.

셔플이 디스크에 쓴 바이트

3단계 앱(spk-shuffle-8)의 가장 최근 로그에서 모든 태스크의 Shuffle Write Metrics.Shuffle Bytes Written 을 더해 /root/spk/shuffle/out/shuffle_bytes.json{"app": "spk-shuffle-8", "shuffle_write_bytes": 정수} 로 쓰세요.

셔플 쓰기는 맵 쪽 태스크가 결과를 파티션별로 나눠 로컬 디스크에 쓰는 일입니다. 부분 집계(HashAggregate 의 partial) 덕에 고객 2만 명 × 맵 태스크 수만큼의 줄만 넘어가서 원본보다 훨씬 작습니다. 이 숫자가 셔플의 실제 비용입니다.

좁은 변환만 있으면 스테이지가 하나다

/root/spk/shuffle/narrow.py 를 앱 이름 spk-shuffle-narrow 로 만들어 결제 완료 주문에 day = to_date(order_ts) 를 더하고 order_id,customer_id,qty,day 만 골라 /root/spk/shuffle/out/narrow 에 Parquet 으로 쓰세요. 이 앱에는 셔플이 하나도 없어야 합니다.

where·withColumn·select 는 줄 하나를 보고 줄 하나를 냅니다. 다른 태스크의 자료가 필요 없으니 한 스테이지 안에서 이어서 돕니다(파이프라이닝). 여기에 orderBydistinct 를 하나만 넣어도 셔플이 생깁니다.

파티션 수를 정한 근거 남기기

/root/spk/shuffle/report.md## 200 개의 파티션 ## AQE 가 합친 것 ## repartition 과 coalesce 세 절을 쓰세요. 첫 절에 2단계의 파일 수, 둘째 절에 4단계의 태스크 수, 셋째 절에 5단계의 두 파일 수를 숫자로 넣으세요.

이 자료 크기라면 셔플 파티션을 몇으로 두겠는지, 그리고 AQE 가 있으니 굳이 손으로 정할 필요가 있는지를 한 줄씩 덧붙이세요. 운영에서 이 결정의 근거가 되는 것은 결과 파일 수와 셔플 바이트입니다.