LabHub
배우기 러닝패스 코스

Apache Spark — The answer to a slow job is in the plan and the event log

Reduce the files and tasks produced by 200 shuffle partitions

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 가 있으니 굳이 손으로 정할 필요가 있는지를 한 줄씩 덧붙이세요. 운영에서 이 결정의 근거가 되는 것은 결과 파일 수와 셔플 바이트입니다.