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

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

explain は Spark が何を読まずに済ませたかを示す領収書だ

TT Labで続きを見る

一言でいうと

Sparkは書いたコードをそのまま実行しません。オプティマイザーが、条件をファイルリーダー側に押し下げ、使わないカラムを切り落とし、不要なパーティションディレクトリをスキップするようプランを書き換えます。その結果はexplainの物理プランにPushedFilters・ReadSchema・PartitionFiltersとして表示され、実行中にAQEがもう一度書き換えた最終プランは、アクションが終わったあとにしか見えません。

なぜプランを読む必要があるのか

同じ結果を出す2つのクエリのうち、片方は3秒、もう片方は3分かかります。コードはほとんど同じに見えます。差は多くの場合、どれだけ読んだかで生じます。条件を1つ関数で包んだせいでParquetファイル全体を読んだとか、select("*")の1行のせいで不要なカラム20個をディスクから読み込んだ、といった具合です。

こうした差は、コードをいくら読んでも見えません。Sparkがコードをどんなプランに変換したかを見る必要があります。その領収書がexplainです。ここで重要な姿勢が1つあります。最適化を信じず、確認します。オプティマイザーはできるときだけ最適化を行い、できなかったときに警告を出しません。

どう動くのか

SQL EXPLAINリファレンスによると、EXTENDEDは4種類のプランを表示します。クエリから取り出したパース済みの論理プラン、名前と型を確定した分析済みの論理プラン、最適化ルールを通った最適化済みの論理プラン、そして実際に実行される物理プランです。PySparkのexplainは、これをモードで選びます。デフォルト(simple)は物理プランだけ、extendedは論理と物理の両方、codegenは生成されたコード、costは統計があるとき論理プランと統計、formattedは物理プランの概要とノードごとの詳細に分けて表示します。

最適化済みの論理プランと物理プランを比べると、オプティマイザーがしたことが見えます。ファイルを読むノードの1行に、ほとんどが集まっています。

+- FileScan parquet [order_id#0,qty#3,day#7]
     PartitionFilters: [isnotnull(day#7), (day#7 = 2026-01-03)]
     PushedFilters: [IsNotNull(qty), GreaterThanOrEqual(qty,5)]
     ReadSchema: struct<order_id:string,qty:int>

条件のプッシュダウン(PushedFilters): where(qty >= 5)がFilterノードとしてだけ存在するのではなく、ファイルリーダーに渡されたという意味です。Parquetの設定で、spark.sql.parquet.filterPushdownのデフォルトはtrueです。渡した条件で何をスキップできるかは、ファイルが持つ統計次第です。統計の活用の節は、Sparkがデータソースから直接読む統計の例として、Parquetメタデータに含まれる件数と最小値・最大値を挙げています。最小値・最大値を見るだけで、条件を満たせないかたまりは開かなくて済みます。

カラムプルーニング(ReadSchema): クエリが最後まで使うカラムだけを読みます。Parquetはカラムごとにまとめて保存するので、2つのカラムだけ読めば、残りのカラムのバイトはディスクから読み込まれません。CSVは行単位なので、この利点はずっと小さくなります。ファイルを一度Parquetに変換しておく理由がここにあります。

パーティションプルーニング(PartitionFilters): day=2026-01-03/のように、値がディレクトリ名に入っているカラムにかかる条件です。Parquetのパーティション検出は、このようなパスからパーティションカラムを自動的に取り出してスキーマに入れます。そのカラムにかかる条件は、ファイルを開く前のディレクトリ一覧の段階で処理されるため、該当しない日付はコストが0です。

連続フィルターの統合: where(a).where(b)のように条件を複数回に分けて書いても、最適化済みの論理プランではFilter 1つに統合されます。そのため、可読性のために条件を分けて書いても、性能は損なわれません。

プッシュダウンが効かなくなる瞬間

プッシュダウンは、条件がカラムそのものと定数の比較であるときに最もうまくいきます。リーダーが理解できるのは、「qtyが5以上」のような単純な形です。upper(status) = 'PAID'のようにカラムを関数で包むと、リーダーに渡せる形ではなくなり、PushedFiltersからその条件が抜けます。結果は同じで、エラーも警告もありません。読む量だけが増えます。このラボで、まさにこの違いをプランで確認します。

パーティションカラムも同じです。日付パーティションにday = '2026-01-03'をかければPartitionFiltersに入りますが、日付文字列を加工して絞り込むと、プルーニングが効かないことがあります。絞り込む条件は、保存されている形のままのカラムにかけるという習慣は、ここから生まれます。

AQE(実行の途中で変わるプラン)

パフォーマンスチューニングのドキュメントによると、適応型クエリ実行(AQE)はランタイム統計で実行の途中にプランを再最適化する手法で、Spark 3.2.0からデフォルトでオンになっています。このことは、プランを読むうえで重要な結果をもたらします。アクションの前に出力したexplainは、最終プランではありません。

アクションの前に出力すると、一番上にAdaptiveSparkPlan isFinalPlan=falseが見えます。シャッフルの段階が実際に終わってサイズがわかってから、AQEがシャッフルパーティションを統合したり、結合の方式を変えたりします。アクションを1回呼んだあとに同じDataFrameのプランをもう一度見ると、isFinalPlan=trueとともにAQEShuffleReadのようなノードが現れます。実行されたプランを確認するには、この最終プランを見る必要があります。Web UIのドキュメントのSQLタブでも、クエリごとにDetailsで4つのプランを展開し、演算子ごとに何行が通過したかといったメトリクスを見られます。

現場での姿

「フィルターをかけたのに、なぜ全部読むのか」: 最もよくある原因は、関数で包んだ条件、型が違う比較(文字列カラムと数値定数)、そしてCSVをそのまま読むことです。3つとも、プランのPushedFilters・ReadSchemaの1行に現れます。

「selectは後でやるから関係ない」という考えは、おおむね正しいです。オプティマイザーが、最後まで使うカラムだけを残すからです。ただし、途中でcache()をしたり、Python UDFに行全体を渡したりすると、その時点で必要なカラムが増えます。常にReadSchemaで確認します。

統計がないとプランが揺らぎます。統計の活用の節は、統計がないか不正確だと、Sparkが良いプランを選べないと書いており、explain(mode="cost")で推定値を、SQL UIのisRuntime=trueで実行中の統計を見るよう案内しています。

実務で本当に大切なこと

次のラボですること

注文CSVを、日付カラムを加えたParquetに変換したあと、explainのextendedとformattedのプランを出力します。条件とカラムの選択がPushedFiltersとReadSchemaにどう表示されるかを確認し、同じ条件をupper()で包んだときにプッシュダウンが消えることを比べます。日付で分けて書いたParquetから1日を選んでPartitionFiltersに条件が表示されるのを見て、顧客別の集計をアクションとして実行したあと、AQEの最終プランとイベントログを読んで、シャッフルパーティション200個のうち実際にシャッフルを読んだタスクがいくつだったかを数えます。最後に、連続したフィルターが1つに統合され、1 + 2のような定数式が事前に計算される様子を、最適化済みの論理プランで確認します。