TT Lab
はじめる
学ぶ 学習パス コース

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

ボット一体がクリックの40%を占めるデータで偏りを測り、解消する

TT Labで続きを見る

目標

クリックの40%をボット1人が作ったデータをユーザーテーブルと結合し、シャッフル後のタスク1つにレコードが集中する様子をイベントログで測ります。AQEの偏り結合・ソルティング・ホットキーの切り出しの3つで解決し、同じホットキーがcount集計ではなぜ問題にならないのかも確認します。

なぜ重要なのか

「ジョブが99%で止まる」という報告の大半は、偏りです。シャッフルはキーのハッシュでパーティションを決めるので、同じキーは必ず1つのタスクに行きます。キー1つが全体の40%なら、タスク1つが40%を抱え込み、残りのタスクがすべて終わったあとも、その1つだけが動き続けます。コアを増やしても意味がありません。そのタスクは分割されないからです。 偏りは、時間よりも分布で見るほうが正確です。ステージの中でタスクごとの読み取りレコード数の最大値と中央値を比べれば、マシンが速くても遅くても同じ数値が出ます。このラボも、時間ではなくその比率で判定します。 対処法は3つです。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を付け、ユーザー側のボットの行を8セットに増やして、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つの節を書いてください(見出しは韓国語で、順に「偏りの形」「3つの対処法」「部分集計」を意味します)。最初の節にステップ3の2つの数値を、3つ目の節にステップ7の2つの数値を、入れてください。

参考

ホットキーを探す

/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行ずつ書いてください。本番でどれを先に試すかも書くとよいでしょう。