Apache Spark — 遅いジョブの答えは実行計画とイベントログにある
シャッフルは最も高くつく一歩で、パーティション 200 はデータを知らずに決めた数だ
一言でいうと
シャッフルは、同じキーを1か所に集めるためにすべてのパーティションがすべてのパーティションにデータを送る処理で、ディスク書き込み・シリアライズ・ネットワークをまとめて支払います。シャッフル後のパーティション数はデフォルトで200で、データサイズとは無関係に決まっている数字です。AQEが実行中に小さなパーティションを統合してくれますが、何が統合されたかと、repartition・coalesceの違いは、自分で確認する必要があります。
なぜシャッフルが問題なのか
filterやwithColumnは、パーティション1つを受け取ってパーティション1つを出します。隣のパーティションを見る必要がないので、タスク1つの中で完結します。こうしたものをナロートランスフォーメーションと呼びます。groupBy・join・distinct・orderByは違います。同じキーがどのパーティションに散らばっているかわからないので、計算するにはすべてのパーティションを調べて、同じキー同士を集める必要があります。
RDDガイドのシャッフルの節は、これをパーティション間でデータを再分配するSparkの仕組みと呼び、コストが大きい理由としてディスクI/O・データのシリアライズ・ネットワークI/Oの3つを挙げています。シャッフルを準備する側はマップタスク、集めて計算する側はリデュースタスクという名前を使いますが、ドキュメントは、この名前はMapReduceに由来するだけで、Sparkのmap・reduce操作とは直接関係がないと付け加えています。
動作は次のとおりです。マップ側のタスクは、結果をメモリに集め、あふれたら宛先パーティションの順にソートしてファイルに書きます。リデュース側のタスクは、すべてのマップタスクが書いたファイルから、自分の担当分のブロックだけを選んで読みます。ドキュメントは、シャッフルがディスクに中間ファイルを大量に作り、リネージを再計算するときに備えて、そのファイルを参照がなくなるまで残しておくと書いています。そのため、シャッフルが多い長いジョブは、ディスクをかなり消費します。
どう動くのか(200という数字)
シャッフル後のパーティションをいくつに分けるかは、誰かが決める必要があります。DataFrameとSQLでは、その値がパフォーマンスチューニングのドキュメントのspark.sql.shuffle.partitionsで、結合や集計でシャッフルするときに使うパーティション数で、デフォルトは200です。
200は、データを見て選んだ数字ではありません。1TBにとっては少なすぎて、パーティション1つが5GB前後になり、メモリがあふれてディスクに流れます。30万行にとっては多すぎて、1つのパーティションに数百行ずつ入り、タスク200個を起動して回収する固定コストが、仕事より大きくなります。その結果をそのままファイルに書くと、書き込みタスクごとにファイルを書くので、小さなファイルが大量にできます。
spark.conf.set("spark.sql.adaptive.enabled", "false") # AQE 를 끄고 날것을 본다
by_day = orders.groupBy("day").agg(F.sum("qty"))
by_day.write.mode("overwrite").parquet("/root/spk/out/by_day") # 셔플 뒤 태스크 200개
spark.conf.set("spark.sql.shuffle.partitions", "8") # 자료에 맞춰 줄인다
AQEのパーティション統合
手で数字を合わせるのが難しい理由は、シャッフル後のサイズが、シャッフルする前にはわからないことにあります。シャッフル後パーティションの統合の節が、この問題を解決します。AQEとspark.sql.adaptive.coalescePartitions.enabledが両方オンのとき(どちらもデフォルトtrue)、マップ側の出力統計を見て隣り合う小さなシャッフルパーティションを統合します。開始パーティション数initialPartitionNumを指定しなければ、spark.sql.shuffle.partitionsと同じです。
どこまで統合するかには詳細があります。目標サイズadvisoryPartitionSizeInBytesのデフォルトは64MBですが、parallelismFirstがデフォルトtrueなので、この目標サイズを無視して、最小サイズminPartitionSize(デフォルト1MB)だけを守りながら並列性を最大にします。ドキュメントは、忙しいクラスターでは小さなタスクが増えないよう、この値をfalseにするよう勧めています。つまり、デフォルト設定のAQEは「64MBにまとめる」のではありません。統合された結果は、最終プランにAQEShuffleRead ... coalescedとして表示されます。
repartitionとcoalesce
パーティション数を直接変える方法は2つあり、コストがまったく違います。repartitionは、ハッシュパーティショニングされた新しいDataFrameを作ります。すべての行を再分配するので、シャッフルが発生します。その代わり、パーティションのサイズが均等になり、カラムを指定すれば、同じ値が同じパーティションに集まります。
coalesceはナロー依存です。ドキュメントの例で1000個から100個に減らすと、シャッフルなしに新しいパーティション1つが既存のパーティション10個を担当します。増やすよう指定しても、現在の数のままです。ただしドキュメントは、coalesce(1)のような急激な縮小が、計算そのものを少ないノードで動かすことになりかねないと警告しています。シャッフルがないので、前の段階まで1つのタスクに統合されるからです。そのような場合は、シャッフルを1回支払ってでも、repartitionのほうが前の段階を並列に動かせます。
RDDガイドは、シャッフルを起こしうる操作にcoalesceも含めていますが、これはRDDのcoalesceがシャッフルするかどうかを引数で受け取るためです。DataFrameのcoalesceは、上の説明のとおりシャッフルがありません。SQLではパーティションヒントのCOALESCE・REPARTITION・REBALANCEで同じことができ、ドキュメントもこれを出力ファイル数を減らすツールとして紹介しています。
現場での姿
小さなデータに200個のファイル: 開発用データで実行したジョブが、結果フォルダーに数十バイトのファイルを大量に作ります。設定を1つ見るだけで原因がわかります。AQEがオンなら大きく減りますが、書き込みの直前にcoalesceでファイル数を決めるほうが確実です。
シャッフルの量はログとUIにあります: Web UIのStagesタブは、ステージごとにShuffle readとShuffle writeをバイトとレコードで表示し、イベントログのタスク終了イベントにも同じメトリクスが入っています。「シャッフルが大きい」という感覚ではなく、数値で語ります。
ナロートランスフォーメーションだけのジョブはステージが1つです: 読んで、絞り込んで、書くだけのジョブは、シャッフルがないのでステージが1つです。ステージが2つ以上あれば、どこかにシャッフルがあるという証拠です。
実務で本当に大切なこと
- シャッフルは、ディスク・シリアライズ・ネットワークをまとめて支払います。ステージの境界がそのままシャッフルです。
- shuffle.partitionsの200は、データを知らないデフォルトです。データサイズに合わせて決めるか、AQEに任せたうえで、結果を確認します。
- デフォルトのAQEは、64MBにまとめません。parallelismFirstがtrueなので、並列性を優先して守ります。
- 減らすときはcoalesce、均等に分けるときはrepartitionです。前者はシャッフルがなく、後者はシャッフルがあります。
- coalesce(1)は、前の段階まで1つのタスクにしてしまいます。大きなデータではrepartitionを使います。
- シャッフルのバイト数は、UIとイベントログで数値として読みます。
次のラボですること
AQEをオフにしたまま顧客別の売上を集計して書き込み、シャッフルパーティション200個が作った結果ファイルの数を数えたあと、シャッフルパーティションを8に減らして違いを見ます。設定をデフォルトに戻して、AQEが統合したあとに実際にシャッフルを読んだタスクの数をイベントログで数え、クリックデータにrepartitionとcoalesceをそれぞれかけて、結果ファイルの数を比べます。シャッフルで書き込んだバイト数をログから直接合計し、ナロートランスフォーメーションだけのジョブにはシャッフルがないことを確認して、レポートにまとめます。