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

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

結合は行の突き合わせ方でコストが十倍変わる

TT Labで続きを見る

一言でいうと

Sparkは、同じ結合を3つの方式(ブロードキャストハッシュ・ソートマージ・シャッフルハッシュ)で行うことができ、どれを選ぶかは片方がどれだけ小さいかによって決まります。その判断は、プランを立てるときの推定と、実行中の実測(AQE)の2回行われます。

なぜ結合戦略を知る必要があるのか

結合は、2つのテーブルからキーが同じ行を1か所に集める処理です。問題は、その行たちが最初は別々のパーティション、別々のエグゼキューターに散らばっていることです。対応する行同士を合わせるには、誰かが移動しなければなりません。何をどれだけ移動させるかが、結合戦略のすべてです。

注文1,000万件と商品300個を結合するとします。注文をキーで再分割して移動させると、1,000万件がネットワークとディスクを通ります。商品300個を全タスクに1セットずつ配れば、注文は元の場所にあり、コピーされるのは300個だけです。結果はまったく同じなのに、移動する量は数万倍違います。実務で「結合が遅い」と言われるものの半分は、この選択が間違っている場合です。

どう動くのか

3つの結合戦略を並べた図です。ブロードキャストハッシュ結合は、小さなテーブル1セットをすべてのタスクにコピーし、大きなテーブルは移動しません。ソートマージ結合は、両側をキーのハッシュでシャッフルしたあと、パーティションごとにソートして並べて走査します。シャッフルハッシュ結合は、両側をシャッフルしたあと、パーティションごとに小さい側でハッシュテーブルを作ります

ブロードキャストハッシュ結合(BroadcastHashJoin): 小さい側のテーブルをドライバーが集めて、すべてのエグゼキューターに1セットずつ送ります。エグゼキューターはそれでハッシュテーブルを作り、大きい側の各パーティションは、その場でハッシュテーブルを引きます。大きい側にシャッフルがありません。そのため最も速いですが、小さい側が本当に小さい必要があります。ドライバーとすべてのエグゼキューターのメモリに丸ごと載るからです。

Sparkがこれを自動的に選ぶ基準が、パフォーマンスチューニングのドキュメントのspark.sql.autoBroadcastJoinThresholdです。デフォルトは10MBで、統計から見たテーブルサイズがこれより小さければブロードキャストします。-1にすると、自動ブロードキャストがオフになります。同じドキュメントのspark.sql.broadcastTimeoutは、ブロードキャストを待つ時間で、デフォルトは300秒です。

ソートマージ結合(SortMergeJoin): 両側を結合キーのハッシュでシャッフルして、同じキーが同じ番号のパーティションに集まるようにします。そのあとパーティションごとに両側をキーでソートし、2本の列を並べて走査しながら対応する行を合わせます。両側ともシャッフルとソートのコストを支払いますが、ソートはメモリが足りなければディスクに流せるので、どんなサイズでも最後まで完了します。2つのテーブルがどちらも大きいときのデフォルトの選択が、これである理由です。

シャッフルハッシュ結合(ShuffledHashJoin): シャッフルはソートマージと同じです。違うのは、パーティションごとにソートの代わりに小さい側でハッシュテーブルを作ることです。ソートを省くので速いことがありますが、パーティション1つ分の小さい側がメモリに収まる必要があります。そのためSparkは、条件が合うときにだけこれを選びます。

from pyspark.sql import functions as F

orders.join(products, "product_id").explain()            # 작으면 BroadcastHashJoin
orders.join(F.broadcast(products), "product_id")         # 크기와 상관없이 브로드캐스트
orders.join(products.hint("shuffle_hash"), "product_id") # ShuffledHashJoin 요청
orders.join(products.hint("merge"), "product_id")        # SortMergeJoin 요청

プランを読むときは、結合の名前よりもその下に何が付いているかを見ます。ソートマージなら、両側の枝にExchange hashpartitioning(product_id, …)とSortが1つずつ付きます。シャッフル2回とソート2回です。ブロードキャストなら、小さい側にだけBroadcastExchangeが付き、大きい側の枝にはExchangeがありません。シャッフルハッシュなら、Exchangeは両側にあるのにSortがありません。この3つの形に目が慣れれば、結合の名前を探さなくても、何を移動したかが見えます。

ヒントとその限界

推定が間違っていることがあります。CSVのように統計がない元データや、フィルターを何度も通った結果は、サイズを推測しにくくなります。そのとき、人間が戦略を指示するのがヒントです。ヒントのドキュメントが挙げる結合ヒントは、BROADCAST(別名BROADCASTJOIN・MAPJOIN)、MERGE(別名SHUFFLE_MERGE・MERGEJOIN)、SHUFFLE_HASH、SHUFFLE_REPLICATE_NLの4つです。BROADCASTヒントが付いた側は、しきい値とは無関係にブロードキャストされます。DataFrame APIでは、broadcast()関数が同じことをします。

両側に異なるヒントを付けると、BROADCAST、MERGE、SHUFFLE_HASH、SHUFFLE_REPLICATE_NLの順に、前のものが優先されます。そして、パフォーマンスチューニングのドキュメントは、はっきり書いています。ヒントは保証ではありません。戦略によっては、特定の結合の種類をサポートしないからです。たとえば左外部結合では、左のすべての行を残す必要があるので、左側をブロードキャストしてハッシュテーブルとして使うことはできません。そのため、ヒントを入れたあとには、必ずプランで実際に何が選ばれたかを確認します。

実行中に変わるプラン(AQE)

適応型クエリ実行(AQE)は、3.2.0からデフォルトでオンです。AQEは、シャッフルが終わったあとに実際に何バイト出たかを見て、残りのプランを組み直します。結合で重要なのは、ソートマージをブロードキャストハッシュに切り替えるルールです。パフォーマンスチューニングのドキュメントによると、実行中の統計で見た片方が、適応型しきい値(spark.sql.adaptive.autoBroadcastJoinThreshold、デフォルトは自動しきい値と同じ)より小さければ切り替えます。

ドキュメントは、ここに正直な但し書きを付けています。最初からブロードキャストで計画したものほど効率的ではありません。すでにシャッフルは起きているからです。その代わり、両側のソートを避けられ、ローカルシャッフル読み込みがオンなら、シャッフルファイルをネットワークを使わずにその場で読みます。プランでは、最初はAdaptiveSparkPlan isFinalPlan=falseだったものが、実行後はisFinalPlan=trueになり、その中で結合の名前が変わったのが見えます。

ソートマージをシャッフルハッシュに切り替えるルールもあります。シャッフル後のすべてのパーティションがspark.sql.adaptive.maxShuffledHashJoinLocalMapThresholdより小さければ切り替えますが、この値のデフォルトが0なので、別にオンにしなければ起こりません。

現場での姿

1つ目は、結合のあとで行数が増えることです。結合キーが片方で一意でなければ、対応する行が掛け算になります。注文1件に対してプロモーションテーブルの同じ商品の行が3つあれば、その注文は3行になります。売上合計が膨らんでレポートが間違うのに、エラーは1つも出ません。結合の前に、小さい側のキーの一意性を数えてみる習慣が、この事故を防ぎます。

2つ目は、ブロードキャストがドライバーを落とすことです。しきい値を大きく上げたり、ヒントを乱発したりすると、数百MBのテーブルがドライバーに集まってから、すべてのエグゼキューターにコピーされます。ドライバーのメモリ不足やブロードキャストのタイムアウトが、このとき起きます。

3つ目は、「存在しないもの」を探すとき、外部結合のあとにnullを除外する代わりにアンチ結合を使うことです。結合のドキュメントは、アンチ結合を、右側に対応する行がない左側の行を返す結合として定義しています。注文が一度もない顧客を探す質問がまさにこの形で、結果には左側のカラムだけが残ります。

実務で本当に大切なこと

次のラボですること

注文と小さな商品テーブルを結合して、Sparkが自分でブロードキャストを選ぶことをプランで確認します。しきい値をオフにしてソートマージに変わるのを見て、broadcastとshuffle_hashのヒントで戦略を直接変えてみます。しきい値を下げたままAQEが実行中にソートマージをブロードキャストに切り替える瞬間を、最終プランで捉え、プロモーションテーブルの重複キーが行を増やす様子を、left_semi結合の行数と比べて数えたあと、アンチ結合で注文のない顧客を探します。