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

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

タスク一つだけが終わらない理由はキー一つにある

TT Labで続きを見る

一言でいうと

シャッフルは同じキーを同じタスクに送るので、キー1つがデータの大きな割合を占めると、そのキーを受け取ったタスク1つがステージ全体を引き止めます。これが偏り(skew)で、解決する方法は、AQEの偏り結合、ソルティング、ホットキーの切り出しの3つです。

なぜタスク1つが問題なのか

ステージは最も遅いタスクが終わって初めて終わります。タスク199個が2秒で終わっても、1つが3分かかればステージは3分です。その間、残りのコアは遊んでいます。エグゼキューターを2倍に増やしても意味がありません。遅いタスク1つは、やはりコア1つで動くからです。

なぜ1つだけ遅いのか。RDDプログラミングガイドは、シャッフルを、1つのキーの値をすべて1か所に集めるために、すべてのパーティションを読むall-to-all操作として説明しています。結合でも集計でも、キーのハッシュで宛先パーティションを決めるので、同じキーは必ず同じタスクに行きます。これがシャッフルの正しさを保証するルールであると同時に、偏りの原因でもあります。クリック24万件のうち、ボット1人が9万6千件を作ったなら、そのユーザーIDを受け取ったパーティション1つに40%が集中します。パーティション数を増やしても解決しません。キー1つは分割されないからです。

どう見つけるのか

偏りは平均では見えません。ステージ全体の読み取りバイトは問題なく見えます。見るべきなのはタスクごとの分布です。Web UIのドキュメントが説明するステージ詳細画面には、すべてのタスクのサマリーメトリクスがあり、その中のDurationとShuffle Read Size / Recordsが、最小・中央値・最大で表示されます。最大が中央値の何倍かが、偏りの大きさです。数十倍なら、タスク1つが仕事をほとんど1人でこなしているという意味です。

イベントログにも同じ数値があります。タスクが終わるたびに残る記録に、そのタスクがシャッフルで読んだレコード数が入っているので、ログ1ファイルがあれば、UIがなくても最大と中央値を直接数えられます。

1つ目の方法(AQEの偏り結合)

パフォーマンスチューニングのドキュメントの偏り結合の最適化は、ソートマージ結合の偏ったパーティションを、同じくらいのサイズの複数のタスクに分割し、反対側の対応するパーティションを必要なだけ複製します。spark.sql.adaptive.enabledとspark.sql.adaptive.skewJoin.enabledが両方オンである必要があり、どちらもデフォルトはtrueです。

どのパーティションが偏っているかは、2つの条件を両方満たしたときに決まります。サイズが中央値のskewedPartitionFactor倍(デフォルト5.0)より大きく、同時にskewedPartitionThresholdInBytes(デフォルト256MB)より大きい必要があります。2つ目の条件のため、小さなラボのデータでは、いくら偏っていても、デフォルト設定では何も起きません。ラボでしきい値を下げる理由です。分割するときに目標にするサイズは、advisoryPartitionSizeInBytes(デフォルト64MB)です。

動作すると、最終プランにSortMergeJoin(skew=true)が現れ、その下のシャッフル読み込みがAQEShuffleRead … coalesced and skewedに変わります。1つ注意点があります。分割によってシャッフルがもう1つ必要になる形の場合、AQEはデフォルトでは適用しません。それでも行いたいときにオンにするのが、spark.sql.adaptive.forceOptimizeSkewedJoin(デフォルトfalse)です。

2つ目の方法(ソルティング)

AQEがない場合や、結合以外の場所の偏りは、人間が解決します。ソルティングは、大きい側のキーに0からN-1までのランダムな数字(ソルト)を付けて、ホットキー1つをN個の異なるキーにします。すると、ハッシュがそれらをN個のパーティションに散らします。その代わり、小さい側は、対応する行を失わないように、すべてのソルト値の数だけ複製する必要があります。

from pyspark.sql import functions as F
N = 8
clicks_s = clicks.withColumn("salt", (F.rand(7) * N).cast("int"))
users_s = users.crossJoin(spark.range(N).withColumnRenamed("id", "salt"))
joined = clicks_s.join(users_s, ["user_id", "salt"]).drop("salt")

代償は明確です。小さい側がN倍になります。そのためNは、偏りを解消できるだけの大きさにとどめ、小さい側が本当に小さいときに使います。

3つ目の方法(ホットキーの切り出し)

ホットキーが数個に決まっているなら、もっと単純な方法があります。そのキーの行だけを抽出して別に処理し、残りは通常どおり結合してから、2つをunionでつなぎます。ホットキー側は、反対側のテーブルでそのキーに該当する行が数行だけなので、ブロードキャストで結合すればシャッフルがまったくありません。残りは、偏りがなくなった均等なデータです。ホットキーを見つけるには、キーごとの件数を数えて上位数個を見れば十分です。

3つのうちどれを選ぶのか

順序は、たいてい次のとおりです。まず、AQEがすでに解決しているかを最終プランで確認します。ソートマージ結合で、偏ったパーティションが2つの条件を超えていれば、設定を何も変えなくても解決されます。AQEが分割する形を思い浮かべると、限界も見えます。偏った側のパーティションを複数の断片に分け、断片ごとに反対側の対応するパーティションを丸ごと付けます。反対側の対応するパーティションも大きければ複製のコストが大きくなり、両側が同じキーで一緒に偏っていれば、断片に分けても対応する相手の数自体は減りません。

AQEで解決しなければ、ホットキーが何個に決まっているかを見ます。ボットアカウント1つ、大口顧客3社のように名前を挙げられるなら、切り出しが最も単純で、結果の説明も簡単です。ホットキーが日ごとに変わる、または数十個あるなら、ソルティングのほうが優れています。どのキーがホットかわからなくても、すべてのキーを均等に散らすからです。どの方法でも、終わったらタスクごとの最大と中央値を測り直して、比率が実際に下がったかを確認します。

偏りが隠れる場合(部分集計)

同じボットのデータでgroupBy("user_id").count()を実行すると、不思議なことに偏りがほとんど見えません。プランを見ると理由があります。シャッフルの前にHashAggregate(partial_count)があり、各マップタスクがキーごとに1行の部分合計だけを送るからです。ボットの9万6千件は、シャッフルの前にマップタスクの数だけの数値に減っています。RDDガイドが、キーごとの合計や平均にはgroupByKeyの代わりにreduceByKey・aggregateByKeyを使うよう勧めているのも、同じ理由です。

そのため、偏りは結合で、そして部分集計のない処理で現れます。ウィンドウ関数はパーティションキーのすべての行を1つのタスクに集める必要があり、リストを集める集計は、部分結果そのものが元の行と同じ大きさです。集計で問題なかったからといって、結合も問題ないだろうと思い込んではいけません。

現場での姿

1つ目は、進捗バーが99%で止まることです。タスク1つだけが残って、数十分動き続けます。そのタスクのシャッフル読み込みレコードが、ほかのタスクの数十倍なら、偏りです。

2つ目は、nullキーがホットキーになることです。値が欠けたカラムにnullが数百万個たまると、nullもハッシュで見れば1つのキーなので、1つのパーティションに集まります。結合条件では、nullはどんな値とも等しくないため、対応する行ができません。しかし外部結合は、対応する行がない行も結果に残す必要があるので、それらの行がそのままシャッフルされます。nullキーの行を先に切り離しておき、結合のあとにunionでつなげば、結果は同じで偏りはなくなります。

3つ目は、偏りは育つことです。ボットや大口顧客が現れた日から、昨日まで問題なかったジョブが遅くなります。コードはそのままなので、データの分布を見なければ原因が見つかりません。結合キー別の上位数個の行数を毎日記録しておけば、ホットキーが育つ兆しを、ジョブが遅くなる前に捉えられます。

実務で本当に大切なこと

次のラボですること

ボット1人がクリックの40%を作ったデータで、ホットなユーザーをまず見つけます。ブロードキャストとAQEをオフにしてソートマージで結合したあと、イベントログからシャッフルを読んだタスクごとのレコード数を取り出し、最大値と中央値を計算します。次に、AQEの偏り結合をオンにし、小さなラボのデータに合わせて偏りのしきい値を下げて、最終プランに偏りの印が現れることを確認します。ソルティングとホットキーの切り出しで同じ結合をもう一度解き、3つの対処法の答えが同じかを比べ、最後に、ユーザー別のgroupBy集計では、部分集計のおかげで偏りがほとんどないことを、同じ2つの数値で確認します。