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

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

同じ結合を四つの戦略で実行し、実行計画で確かめる

TT Labで続きを見る

目標

注文と商品の結合を、ブロードキャストハッシュ・ソートマージ・シャッフルハッシュの3つの戦略で実行してプランと結果を比べ、AQEが実行中にソートマージをブロードキャストに切り替えることをイベントログで確認します。キーが重複するテーブルとの結合が行を増やすことと、反対側に対応する行がない行を選ぶ結合も行ってみます。

なぜ重要なのか

結合はSparkジョブで最も高価な処理であることが多く、そのコストは戦略が決めます。片方が小さければ、小さい側をすべてのタスクにコピーして(ブロードキャスト)、大きい側をシャッフルせずに結合します。どちらも大きければ、両側をキーでシャッフルしたあとソートして、突き合わせながら結合します(ソートマージ)。シャッフル後に片方をハッシュテーブルにして結合するシャッフルハッシュは、ソートを省く代わりにメモリを使います。 Sparkは統計でサイズを推測して戦略を選びます。小さいと見なす基準がspark.sql.autoBroadcastJoinThreshold(デフォルト10MB)です。推測が間違えば、戦略も間違います。フィルターのあとは実際には小さいのに元のサイズで推測してソートマージを選んだり、逆に大きなテーブルをブロードキャストしてドライバーのメモリがあふれたりします。AQEは、シャッフルが終わったあとの実際のサイズで再判断して、これを直します。 結合結果の行数は、キーの重複が決めます。片方のキーが一意だと信じていたのにそうでなければ、結合はエラーなしに行を掛け算してしまい、後ろの合計がすべて膨らみます。

ステップ

  1. /root/spk/join/common.pyに4つのテーブルを読み込む関数とカテゴリ別売上の関数を置き、/root/spk/join/auto.py(アプリspk-join-auto、デフォルト設定)で、カテゴリ別売上を、/root/spk/join/out/by_categoryにヘッダー行ありのCSV(category・revenue)で書き込んでください。
  2. /root/spk/join/smj.py(アプリspk-join-smj、spark.sql.autoBroadcastJoinThreshold=-1・spark.sql.adaptive.enabled=false)で、同じ結果を、/root/spk/join/out/by_category_smjに書き込んでください。
  3. /root/spk/join/hint.py(アプリspk-join-hint、ステップ2と同じ設定)で、商品側にbroadcastヒントを付けて、/root/spk/join/out/by_category_hintに書き込んでください。
  4. /root/spk/join/shash.py(アプリspk-join-shash、ステップ2と同じ設定)で、商品側にshuffle_hashヒントを付けて、/root/spk/join/out/by_category_shashに書き込んでください。
  5. /root/spk/join/aqe.py(アプリspk-join-aqe、spark.sql.autoBroadcastJoinThreshold=100k、AQEオン)で、決済完了の注文をtier == 'vip'の顧客と結合して、都市別の注文数を、/root/spk/join/out/vip_by_cityにCSV(city・orders)で書き込んでください。
  6. /root/spk/join/dup.py(アプリspk-join-dup)で、決済完了の注文をプロモーションテーブルとそのまま結合した行数と、left_semiで結合した行数を、/root/spk/join/out/dup.jsonに{"naive": 정수, "semi": 정수}の形式で書き込んでください(プレースホルダーは順に整数、整数です)。
  7. /root/spk/join/anti.py(アプリspk-join-anti)で、一度も注文していない顧客のcustomer_idを、/root/spk/join/out/no_ordersにCSVで書き込んでください。
  8. /root/spk/join/report.mdに、## 네 가지 전략・## AQE 의 전환・## 키 중복の3つの節を書いてください(見出しは韓国語で、順に「4つの戦略」「AQEによる切り替え」「キーの重複」を意味します)。3つ目の節には、ステップ6の2つの数値を入れてください。

参考

小さい側は自動的にブロードキャストされる

/root/spk/join/common.pyに、4つのテーブル(注文・商品・顧客・プロモーション)をスキーマを指定して読み込む関数と、カテゴリ別売上(決済完了の注文×商品、category・revenue=qty×priceの合計)の関数を置き、/root/spk/join/auto.pyを、アプリ名spk-join-auto(設定はデフォルト)で作成して、結果を、/root/spk/join/out/by_categoryにヘッダー行ありのCSVで書き込んでください。

商品テーブルは数KBなので、しきい値(10MB)よりはるかに小さいです。Sparkはこれをすべてのタスクにコピーし、注文側はシャッフルしません。採点ツールは、イベントログのプランにBroadcastHashJoinがあるか、カテゴリ別売上が元データと一致するかを確認します。

しきい値をオフにするとソートマージになる

/root/spk/join/smj.pyを、アプリ名spk-join-smj、設定spark.sql.autoBroadcastJoinThreshold=-1・spark.sql.adaptive.enabled=falseで作成し、ステップ1と同じ結果を、/root/spk/join/out/by_category_smjに書き込んでください。

しきい値-1は「自動ブロードキャストなし」です。これで、注文と商品の両側がproduct_idでシャッフルされ、ソートされたあとに結合されます。結果はステップ1と1行も違ってはいけません。AQEをオフにするのは、実行中に戦略が変わらないようにするためです。

ヒントでブロードキャストを強制する

/root/spk/join/hint.pyを、アプリ名spk-join-hint、ステップ2と同じ設定(しきい値-1、AQEオフ)で作成しますが、商品側をF.broadcast(p)で包んで、/root/spk/join/out/by_category_hintに書き込んでください。

ヒントは統計より優先されます。しきい値をオフにしてあっても、ヒントがあればブロードキャストします。逆に言えば、大きなテーブルに気軽に付けたブロードキャストヒントは、ドライバーとすべてのエグゼキューターのメモリをそのまま消費します。

シャッフルハッシュ結合(ソートを省く代わりに)

/root/spk/join/shash.pyを、アプリ名spk-join-shash、ステップ2と同じ設定で作成しますが、商品側にp.hint("shuffle_hash")を付けて、/root/spk/join/out/by_category_shashに書き込んでください。

シャッフルハッシュ結合は、両側をキーでシャッフルしたあと、小さい側のパーティションをハッシュテーブルにして、大きい側を流し込みます。ソートがないので速いことがありますが、ハッシュテーブルがメモリに収まる必要があります。プランからSort演算子が消えていることを確認してください。

AQEが実行中に戦略を変える

/root/spk/join/aqe.pyを、アプリ名spk-join-aqe、設定spark.sql.autoBroadcastJoinThreshold=100k(AQEはオンのまま)で作成し、決済完了の注文をtier == 'vip'の顧客とcustomer_idで結合して、都市別の注文数を、/root/spk/join/out/vip_by_cityにヘッダー行ありのCSV(city・orders)で書き込んでください。

顧客テーブルのファイルは100KBより大きいので、最初のプランはソートマージです。しかし、vipで絞り込んだあとの実際のサイズは数十KBです。AQEは、シャッフルのマップの段階が終わったあと、そのサイズを見て、残りのプランをブロードキャストに切り替えます。採点ツールは、同じ実行の最初のプランと最終プランを比べます。

キーが重複すると、結合が行を増やす

/root/spk/join/dup.pyを、アプリ名spk-join-dupで作成し、決済完了の注文をプロモーションテーブル(/data/shop/promos.csv)とproduct_idでそのまま結合した行数と、left_semiで結合した行数を、/root/spk/join/out/dup.jsonに{"naive": 정수, "semi": 정수}の形式で書き込んでください(プレースホルダーは順に整数、整数です)。

プロモーションテーブルには、コードが2つある商品があります。その商品の注文は、そのまま結合すると2行になります。left_semiは「対応する行があるか」だけを見るので、左側の行を増やしません。売上をプロモーションの有無で分けるとき、どちらを使うべきかを考えてみてください。

対応する行がない側を選ぶ(left_anti)

/root/spk/join/anti.pyを、アプリ名spk-join-antiで作成し、一度も注文していない(状態は問わない)顧客のcustomer_idを、/root/spk/join/out/no_ordersにヘッダー行ありのCSVで書き込んでください。

not inサブクエリや、left joinのあとのnull除外でも実現できますが、left_antiは意味がそのまま見え、nullが混ざったキーでも混乱しません。プランで結合の種類がLeftAntiと出力されるかを見てください。

どの戦略がいつ合うのかを残す

/root/spk/join/report.mdに、## 네 가지 전략・## AQE 의 전환・## 키 중복の3つの節を書いてください(見出しは韓国語で、順に「4つの戦略」「AQEによる切り替え」「キーの重複」を意味します)。最初の節にはステップ1–4で見た結合演算子の名前を、3つ目の節にはステップ6の2つの数値を、入れてください。

最初の節には、各戦略が何をシャッフルして何をメモリに載せるのかを、2つ目の節には、最初のプランと最終プランがどう違ったかを、3つ目の節には、行がいくつ増えたかを書いてください。