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

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

explain で Parquet をどれだけ読まずに済むか確かめる

TT Labで続きを見る

目標

注文をParquetに変換したうえで、同じ質問が条件の書き方によってファイルをどれだけ読まずに済むか(プッシュダウン・カラムプルーニング・パーティションプルーニング)を、実行プランで読みます。適応型実行(AQE)が動いたあとにプランをどう変えたかも、イベントログで確認します。

なぜ重要なのか

Sparkが遅いときに最初に見るのは、コードではなくプランです。同じ結果を出す2行のうち、片方はファイルから必要なカラムと行だけを読み、もう片方はすべて読んでから捨てます。その差は結果には現れず、プランにだけ現れます。 Parquetのような列指向形式は、2つのことができます。必要なカラムだけを読み(カラムプルーニング、プランのReadSchema)、行グループの最小値・最大値の統計で、条件に合わないかたまりをスキップします(条件のプッシュダウン、プランのPushedFilters)。さらに、ディレクトリで分けたパーティションを条件で選べば、ディレクトリごとスキップします(PartitionFilters)。3つとも、オプティマイザーが条件を理解できるときにだけ起こります。カラムを関数で包んだ瞬間に、プッシュダウンは消えます。 AQEは、シャッフルが終わったあとに実際のサイズを見て、残りのプランを書き換えます。そのため、実行前のプラン(isFinalPlan=false)と実行後のプランは異なる場合があり、本当に何が動いたかは、実行後のプランかイベントログで見ます。

ステップ

  1. /root/spk/plan/convert.py(アプリspk-plan-convert)で、注文をスキーマを指定して読み込み、order_date(日付)カラムを加えて、/root/spk/plan/orders_pqにParquetで書き込んでください。
  2. /root/spk/plan/planhelp.py(ヘルパー)にexplain・field・items関数を置き、/root/spk/plan/explain.py(アプリspk-plan-explain)で、返金注文のorder_id・qtyを選ぶクエリのextended・formattedのプランを、/root/spk/plan/out/extended.txt・/root/spk/plan/out/formatted.txtに保存して、そのクエリをcollect()してください。
  3. /root/spk/plan/pushdown.py(アプリspk-plan-push)で、status == 'refunded'かつqty >= 5の注文のorder_id・qty・channelを、/root/spk/plan/out/pushdownにCSVで書き込み、そのプランのPushedFiltersの項目一覧とReadSchemaを、/root/spk/plan/out/pushdown.jsonに書き込んでください。
  4. /root/spk/plan/nopush.py(アプリspk-plan-nopush)で、同じ質問をupper(status) == 'REFUNDED'に変えて、/root/spk/plan/out/nopushに書き込み、そのプランのPushedFiltersの項目一覧を、/root/spk/plan/out/nopush.jsonに書き込んでください。
  5. /root/spk/plan/partition.py(アプリspk-plan-part)で、/root/spk/plan/orders_by_dayにorder_dateで分けて書き込み、2026-02-14の1日を選んで数えた行数とプランのPartitionFiltersを、/root/spk/plan/out/partition.jsonに書き込んでください。
  6. /root/spk/plan/aqe.py(アプリspk-plan-aqe)で、顧客別の注文数をcollect()したあと、プランを、/root/spk/plan/out/aqe_final.txtに保存し、そのアプリのログでシャッフルを読んだタスク数を数えて、/root/spk/plan/out/aqe.jsonに{"shuffle_partitions": 200, "reduce_tasks": 정수}の形式で書き込んでください(プレースホルダーは整数です)。
  7. /root/spk/plan/optimize.py(アプリspk-plan-opt)で、qty > 1とqty > 3の2つのwhereをつなぎ、1 + 2をthreeとして選ぶクエリのextendedプランを、/root/spk/plan/out/optimized.txtに保存して、collect()してください。
  8. /root/spk/plan/report.mdに、## 밀어내기・## 가지치기・## AQEの3つの節を書いてください(最初の2つの見出しは韓国語で、順に「プッシュダウン」「プルーニング」を意味します)。2つ目の節にはステップ5の行数を、3つ目の節にはステップ6のタスク数を、数値で入れてください。

参考

CSVをParquetに変換する

/root/spk/plan/convert.pyを、アプリ名spk-plan-convertで作成し、/data/shop/orders.csvをスキーマを指定して読み込んでorder_date = to_date(order_ts)カラムを加え、/root/spk/plan/orders_pqにParquetで書き込んでください(分割はしません)。

列指向形式のファイルは、カラムごとに別々に保存され、行グループごとに最小値・最大値の統計を持ちます。この2つが、以降のすべてのステップで「読む量を減らす」ための材料です。採点ツールは、行数とカラム一覧、ファイルが本当にParquetか(先頭の4バイトがPAR1)を確認します。

explainのモード(論理プランから物理プランまで)

/root/spk/plan/planhelp.py(ヘルパー)に、explainの出力を文字列として受け取るexplain(df, mode)、formattedプランから이름: 값(プレースホルダーは順に名前、値です)の形の行の値を取り出すfield(plan, name)、[A(x,1), B(y)]を項目の一覧に分けるitems(text)を置いてください。/root/spk/plan/explain.pyを、アプリ名spk-plan-explainで作成し、orders_pqでstatus == 'refunded'の行のorder_id・qtyを選ぶクエリのextended・formattedのプランを、/root/spk/plan/out/extended.txt・/root/spk/plan/out/formatted.txtに保存したあと、そのクエリをcollect()してください。

extendedは、パース済みの論理プラン、分析済みの論理プラン、最適化済みの論理プラン、物理プランの順に4つの節を出力します。分析の段階でカラム名が実際のカラム(#番号)に結び付けられ、最適化の段階で条件がスキャンの近くまで下りてきます。formattedは、その物理プランを演算子ごとに書き出します。採点ツールは、formattedのファイルが実際に動いたプラン(イベントログ)と同じかを確認します。

条件のプッシュダウンとカラムプルーニング

/root/spk/plan/pushdown.pyを、アプリ名spk-plan-pushで作成し、orders_pqでstatus == 'refunded'かつqty >= 5の行のorder_id・qty・channelを、/root/spk/plan/out/pushdownにヘッダー行ありのCSVで書き込み、そのクエリのformattedプランにあるPushedFiltersの項目一覧とReadSchemaの値を、/root/spk/plan/out/pushdown.jsonに{"pushed_filters": [문자열…], "read_schema": 문자열}の形式で書き込んでください(プレースホルダーは順に文字列、文字列です)。

PushedFiltersはParquetの読み込みに渡した条件で、ReadSchemaは実際に読むカラムです。選ばなかったcustomer_id・product_idがReadSchemaから消え、条件にだけ使ったstatusは残ることを確認してください。採点ツールは、自分で書いた一覧が、実際に動いたプランのものと1文字ずつ一致しているかを確認します。

関数で包むとプッシュダウンが消える

/root/spk/plan/nopush.pyを、アプリ名spk-plan-nopushで作成し、ステップ3と同じ質問をF.upper("status") == "REFUNDED"とqty >= 5に変えて、/root/spk/plan/out/nopushにヘッダー行ありのCSVで書き込み、そのプランのPushedFiltersの項目一覧を、/root/spk/plan/out/nopush.jsonに{"pushed_filters": [문자열…]}の形式で書き込んでください(プレースホルダーは文字列です)。

結果の行はステップ3とまったく同じです。変わるのはプランです。Parquetはupper(status)が何かを知らないので、その条件は読み込みに渡されず、ファイルをすべて読んだあとにSparkが絞り込みます。実務では、日付カラムをdate_formatで包んで比較して、この落とし穴にはまることがよくあります。

パーティションプルーニング(ディレクトリごとスキップする)

/root/spk/plan/partition.pyを、アプリ名spk-plan-partで作成し、orders_pqをpartitionBy("order_date")で、/root/spk/plan/orders_by_dayに書き込み、読み直してorder_date == '2026-02-14'の行を数えた値とそのクエリプランのPartitionFiltersの値を、/root/spk/plan/out/partition.jsonに{"rows": 정수, "partition_filters": 문자열}の形式で書き込んでください(プレースホルダーは順に整数、文字列です)。

分割したカラムは、ファイルの中ではなくディレクトリ名(order_date=2026-02-14)に入ります。そのカラムにかける条件はディレクトリを選ぶのに使われ、選ばれなかった日付のファイルは開きもしません。採点ツールは、イベントログの「読んだパーティション数」のメトリクスが1かも確認します。

AQEが実行後に書き換えたプラン

/root/spk/plan/aqe.pyを、アプリ名spk-plan-aqeで作成し、orders_pqの顧客別の注文数(groupBy("customer_id").count())をcollect()したあとに、explain(mode="simple")の出力を、/root/spk/plan/out/aqe_final.txtに保存してください。そして、そのアプリのイベントログでシャッフルを読んだタスク数を数えて、/root/spk/plan/out/aqe.jsonに{"shuffle_partitions": 200, "reduce_tasks": 정수}の形式で書き込んでください(プレースホルダーは整数です)。

spark.sql.shuffle.partitionsは200ですが、このデータのシャッフルは数MBしかないため、AQEがパーティションを統合します。実行後のプランには、isFinalPlan=trueと統合されたシャッフル読み込み(AQEShuffleRead)が見えます。シャッフルを読んだタスク数は、ログのSparkListenerTaskEndでTotal Records Readが0より大きいタスクを数えると求まります。

オプティマイザーが統合し、畳み込むもの

/root/spk/plan/optimize.pyを、アプリ名spk-plan-optで作成し、orders_pq.where(qty > 1).where(qty > 3).select("order_id", (lit(1) + lit(2)).alias("three"))のextendedプランを、/root/spk/plan/out/optimized.txtに保存して、そのクエリをcollect()してください。

分析済みの論理プランにはFilterが2つあり、(1 + 2)がそのまま残っています。最適化済みの論理プランでは、2つの条件が1つのFilterに統合され(フィルターの統合)、1 + 2は3に畳み込まれます(定数畳み込み)。自分が書いたコードの形と、実際に動く形が違うことを確認するステップです。

読まずに済んだ分を数値で残す

/root/spk/plan/report.mdに、## 밀어내기・## 가지치기・## AQEの3つの節を書いてください(最初の2つの見出しは韓国語で、順に「プッシュダウン」「プルーニング」を意味します)。2つ目の節にはステップ5の行数を、3つ目の節にはステップ6のシャッフル読み込みタスク数を、数値で入れてください。

最初の節には、ステップ3・4でPushedFiltersがどう変わったか、2つ目の節には、パーティションプルーニングで何をスキップしたか、3つ目の節には、200からいくつに減ったかを書いてください。