Apache Spark — 遅いジョブの答えは実行計画とイベントログにある
ボット一体がクリックの40%を占めるデータで偏りを測り、解消する
目標
クリックの40%をボット1人が作ったデータをユーザーテーブルと結合し、シャッフル後のタスク1つにレコードが集中する様子をイベントログで測ります。AQEの偏り結合・ソルティング・ホットキーの切り出しの3つで解決し、同じホットキーがcount集計ではなぜ問題にならないのかも確認します。
なぜ重要なのか
「ジョブが99%で止まる」という報告の大半は、偏りです。シャッフルはキーのハッシュでパーティションを決めるので、同じキーは必ず1つのタスクに行きます。キー1つが全体の40%なら、タスク1つが40%を抱え込み、残りのタスクがすべて終わったあとも、その1つだけが動き続けます。コアを増やしても意味がありません。そのタスクは分割されないからです。 偏りは、時間よりも分布で見るほうが正確です。ステージの中でタスクごとの読み取りレコード数の最大値と中央値を比べれば、マシンが速くても遅くても同じ数値が出ます。このラボも、時間ではなくその比率で判定します。 対処法は3つです。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を付け、ユーザー側のボットの行を8セットに増やして、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つの節を書いてください(見出しは韓国語で、順に「偏りの形」「3つの対処法」「部分集計」を意味します)。最初の節にステップ3の2つの数値を、3つ目の節にステップ7の2つの数値を、入れてください。
参考
- 元データ:
/data/clicks/clicks.jsonl(user_id, page, ts, ms。24万行)、/data/clicks/users.csv(user_id, segment。ボットはsegmentがbot)。 - タスクごとの読み取りレコードは、
SparkListenerTaskEndのTask Metrics.Shuffle Read Metrics.Total Records Readです。中央値(低いほう)は、昇順に並べたn個のうち(n-1)//2番目(0から数えて)の値です。 - スクリプトは
/root/spk/skewに置き、そこで実行してください。ログを読む小さなツール(skewtool.py)を作っておくと、ステップ3・7が1行になります。 - よくあるミス: ステップ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位がほかと何倍も差があれば、そのキー1つがタスク1つを引き止めます。
そのまま結合するとタスク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万行あまりのクリックがすべて1つのパーティションに行きます。採点ツールは、そのステージで最大値が中央値の何倍かを確認します。
偏りを数値で見る(最大値と中央値)
ステップ2のアプリ(spk-skew-plain)の最新のログで、シャッフルを読んだステージごとにタスクごとのTotal Records Readを集め、最大値が最も大きいステージを選んで、/root/spk/skew/out/skew.jsonに{"stage_id": 정수, "max_records": 정수, "median_records": 정수}の形式で書き込んでください(プレースホルダーは順に整数、整数、整数です)。中央値は、昇順のn個のうち(n-1)//2番目(0から数えて)の値です。
結合ステージは両側のシャッフルを一緒に読むので、ボットが入ったパーティションのタスクは、ボットのクリックすべてと、ボットユーザーの1行を読みます。最大値÷中央値が、偏りの大きさです。この比率が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を付け、ユーザー側はボットの行だけsalt0–7の8セットに増やしたうえで(残りは0)、user_id・saltで結合した結果を、/root/spk/skew/out/by_segment_saltに書き込んでください。
結合キーにソルトが入るので、ボットのクリックが8つのパーティションに散らばります。反対側のボットの行が8セットあれば、どのソルトに行っても対応する行があります。結果はステップ2とまったく同じでなければならず、採点ツールは、偏りの比率がステップ2の半分未満に下がったかを確認します。
対処法3(ホットキーの切り出し)
/root/spk/skew/split.pyを、アプリ名spk-skew-split、ステップ2と同じ設定で作成し、ボットのクリックはボットユーザーの1行とbroadcast結合し、残りのクリックは通常の結合をしたあと、unionByNameでつなげて、/root/spk/skew/out/by_segment_splitに書き込んでください。
ホットキーが1つで、その名前がわかっているなら、最も単純で確実な対処法です。そのキーの反対側は1行なのでブロードキャストは無料で、残りは均等に分散します。プランに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です。マップタスクごとにユーザー別に先に数えておくので、シャッフルに渡されるのは、マップタスクあたりユーザー1行です。偏りが問題になるのは、このように先に減らせない処理(結合や、collect_listのようなもの)です。
偏りと対処法を数値で残す
/root/spk/skew/report.mdに、## 쏠림의 모양・## 세 가지 처방・## 부분 집계の3つの節を書いてください(見出しは韓国語で、順に「偏りの形」「3つの対処法」「部分集計」を意味します)。最初の節にステップ3の最大値・中央値を、3つ目の節にステップ7の最大値・中央値を、数値で入れてください。
2つ目の節には、3つの対処法がそれぞれ何を変えるのか(コードか設定か、何を複製するのか)を、1行ずつ書いてください。本番でどれを先に試すかも書くとよいでしょう。