Apache Spark — 遅いジョブの答えは実行計画とイベントログにある
シャッフルパーティション 200 が生んだファイルとタスクを減らす
目標
顧客別の売上集計を、AQEをオフにしたデフォルト(シャッフルパーティション200)、手で減らした値(8)、AQEの3通りで実行し、シャッフル後のタスク数と結果ファイル数がどう変わるかをイベントログで測ります。repartitionとcoalesceが、シャッフルの有無とファイル数でどう違うかも確認します。
なぜ重要なのか
groupByやjoinのようなワイドトランスフォーメーションは、同じキーを1つのタスクに集める必要があるので、シャッフルを生みます。前のステージが結果をパーティションごとのファイルに書き、後ろのステージがその断片を集めて読みます。後ろのステージのタスク数はspark.sql.shuffle.partitionsが決め、デフォルトは200です。
200は、大きなクラスターを基準にした数字です。数MBのデータに200を使うと、タスク200個がそれぞれ少しずつ仕事をし、結果をファイルに書くと小さなファイルが200個できます。逆に数百GBに200だと、タスク1つが数GBを抱え込み、ディスクに流れます。そのためこの数字はデータサイズに合わせる必要があり、AQEはシャッフルが終わったあとに実際のサイズを見て、小さな断片を統合することで、この仕事を代わりに行います。
パーティション数を変える方法は2つあります。repartition(n)はシャッフルで断片を新しく作り、coalesce(n)はシャッフルなしで隣り合う断片をくっつけます。coalesceは安価ですが、減らすことしかできません。断片2つを4つにすることはできません。
ステップ
- /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で書き込んでください。 - その結果フォルダーのデータファイル(
part-で始まるもの)の数を数えて、/root/spk/shuffle/out/files_200.txtに整数で記入してください。 - /root/spk/shuffle/eight.py(アプリ
spk-shuffle-8、AQEオフ、spark.sql.shuffle.partitions=8)で、同じ結果を、/root/spk/shuffle/out/by_customer_8に書き込んでください。 - /root/spk/shuffle/aqe.py(アプリ
spk-shuffle-aqe、設定はデフォルト)で、同じ結果を、/root/spk/shuffle/out/by_customer_aqeに書き込み、そのアプリでシャッフルを読んだタスク数を、/root/spk/shuffle/out/aqe_tasks.txtに整数で記入してください。 - クリックの元データを、/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で書き込み、2つのフォルダーのデータファイル数を、/root/spk/shuffle/out/repart.jsonに{"repartition_files": 정수, "coalesce_files": 정수}の形式で書き込んでください(プレースホルダーは順に整数、整数です)。 - ステップ3のアプリがシャッフルで書き込んだバイト数の合計をログから求め、/root/spk/shuffle/out/shuffle_bytes.jsonに
{"app": "spk-shuffle-8", "shuffle_write_bytes": 정수}の形式で書き込んでください(プレースホルダーは整数です)。 - /root/spk/shuffle/narrow.py(アプリ
spk-shuffle-narrow)で、決済完了の注文にdayカラムを加えて、order_id・customer_id・qty・dayだけを、/root/spk/shuffle/out/narrowに書き込んでください。このアプリにはシャッフルがあってはいけません。 - /root/spk/shuffle/report.mdに、
## 200 개의 파티션・## AQE 가 합친 것・## repartition 과 coalesceの3つの節を書いてください(見出しは韓国語で、順に「200個のパーティション」「AQEが統合したもの」「repartitionとcoalesce」を意味します)。各節に、ステップ2・4・5の数値を入れてください。
参考
- スクリプトは
/root/spk/shuffleに置き、そこでspark-submitしてください(from common import …)。 - 設定は
SparkSession.builder.config("키", "값")かspark-submit --conf 키=값で指定します(プレースホルダーはキーと値です)。どちらの場合もイベントログに残ります。 - ログから数値を取り出す小さなツール
logtool.pyを作っておくと、ステップ4・6が1行になります。シャッフルを読んだタスクは、SparkListenerTaskEndのShuffle Read Metrics.Total Records Readが0より大きいタスク、シャッフルで書き込んだバイト数はShuffle Write Metrics.Shuffle Bytes Writtenの合計です。 - よくあるミス: 同じ名前で何度も実行して古いログを数えること(最新のものを数えてください)、
_SUCCESSや.crcまでファイル数に入れること、coalesceでパーティションを増やせると思い込むこと。 - 公式ドキュメント: RDD Programming Guide — Shuffle operations・Performance Tuning — Adaptive Query Execution・Configuration・Monitoring — REST API·metrics
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に整数1つで記入してください。
シャッフル後のタスク1つが、ファイル1つを書きます(空のパーティションはファイルを書かないことがあります)。顧客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と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.jsonlを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で書き込んでください。2つのフォルダーのデータファイル数を、/root/spk/shuffle/out/repart.jsonに{"repartition_files": 정수, "coalesce_files": 정수}の形式で書き込んでください(プレースホルダーは順に整数、整数です)。
18MBの元データが、local[2]で何個の断片として読まれるかを、まず見てください(rdd.getNumPartitions())。repartitionはシャッフルでちょうど4つを作り、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万人×マップタスク数の行だけが渡されるため、元データよりはるかに小さくなります。この数値が、シャッフルの実際のコストです。
ナロートランスフォーメーションだけならステージは1つ
/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で書き込んでください。このアプリには、シャッフルが1つもあってはいけません。
where・withColumn・selectは、1行を見て1行を出します。ほかのタスクのデータが必要ないので、1つのステージの中で続けて動きます(パイプライン化)。ここにorderByやdistinctを1つ入れるだけでも、シャッフルが生じます。
パーティション数を決めた根拠を残す
/root/spk/shuffle/report.mdに、## 200 개의 파티션・## AQE 가 합친 것・## repartition 과 coalesceの3つの節を書いてください(見出しは韓国語で、順に「200個のパーティション」「AQEが統合したもの」「repartitionとcoalesce」を意味します)。最初の節にステップ2のファイル数、2つ目の節にステップ4のタスク数、3つ目の節にステップ5の2つのファイル数を、数値で入れてください。
このデータサイズなら、シャッフルパーティションをいくつにするか、そしてAQEがあるのでわざわざ手で決める必要があるかを、1行ずつ添えてください。本番でこの判断の根拠になるのは、結果ファイル数とシャッフルバイト数です。