Apache Spark — 遅いジョブの答えは実行計画とイベントログにある
月別売上・分類別順位・累積和を SQL と DataFrame で
目標
同じ質問をSpark SQLとDataFrame APIでそれぞれ解いて同じ物理プランが出ることを確認し、ウィンドウ関数で順位・累積合計・前月比を計算します。最後に、正確なユニーク数と近似ユニーク数を比べます。
なぜ重要なのか
SparkにおいてSQLとDataFrameは、2つのエンジンではありません。どちらも同じ論理プランに変換され、同じオプティマイザー(Catalyst)を通り、同じ物理プランで実行されます。そのため、「SQLのほうが速い」や「DataFrameのほうが速い」は、たいてい間違った問いです。チームが何で書いても、プランを見れば同じか違うかがすぐにわかります。 集計の次に最もよく使うのがウィンドウ関数です。行を減らさずに隣の行を見られるため、順位・累積合計・前月比のようなレポートの数値は、すべてここから出てきます。ウィンドウの枠(パーティション・並び順・範囲)を誤って設定すると、エラーなしに間違った数値が出るので、枠を言葉で説明できる必要があります。 ユニーク数を数えるのは、シャッフルが大きい処理です。数パーセントの誤差を受け入れれば、HyperLogLog++ではるかに安く数えられます。ダッシュボードには近似、請求書には正確です。どちらを使うかは、数値の使い道が決めます。
ステップ
- /root/spk/sql/common.pyに、注文と商品をスキーマを指定して読み込み、一時ビュー
orders・productsとして登録するload(spark)を作成してください。/root/spk/sql/sql_monthly.py(アプリspk-sql-monthly)では、SQLで月別売上(カラムはmonthとrevenue)を、/root/spk/sql/out/monthly_sqlにヘッダー行ありのCSVで書き込んでください。 - /root/spk/sql/df_monthly.py(アプリ
spk-sql-df)では、spark.sqlを使わずDataFrame APIで、同じ結果を、/root/spk/sql/out/monthly_dfに書き込んでください。 - /root/spk/sql/plans.py(アプリ
spk-sql-plans)では、2つの方法のexplain(mode="formatted")の出力を、/root/spk/sql/out/plan_sql.txt・/root/spk/sql/out/plan_df.txtに保存したうえで、2つのDataFrameをそれぞれcollect()で実際に実行してください。 - /root/spk/sql/top3.py(アプリ
spk-sql-top3)では、カテゴリ(category)ごとに売上上位3商品を、/root/spk/sql/out/top3にカラムcategory・product_id・revenue・rankで書き込んでください。同点なら、product_idが前のほうが上です。 - /root/spk/sql/running.py(アプリ
spk-sql-running)では、2026年1月のチャネル別の日次売上とチャネル内の累積合計を、/root/spk/sql/out/runningにカラムchannel・day・revenue・runningで書き込んでください。 - /root/spk/sql/mom.py(アプリ
spk-sql-mom)では、チャネル別の月売上と前月売上、増加率(パーセント、小数第2位を四捨五入)を、/root/spk/sql/out/momにカラムchannel・month・revenue・prev_revenue・growth_pctで書き込んでください。最初の月の前月の値は空にしておきます。 - /root/spk/sql/distinct.py(アプリ
spk-sql-distinct)では、決済が完了した注文をした顧客の数を、正確に(countDistinct)、そしてapprox_count_distinct(rsd=0.05)で数えて、/root/spk/sql/out/distinct.jsonに{"exact": 정수, "approx": 정수}の形式で書き込んでください(プレースホルダーは順に整数、整数です)。 - /root/spk/sql/report.mdに、
## SQL 과 DataFrame・## 윈도 함수・## 근사 집계の3つの節を書いてください(見出しは韓国語で、順に「SQLとDataFrame」「ウィンドウ関数」「近似集計」を意味します)。3つ目の節には、ステップ7の2つの数値を入れてください。
参考
- 元データ:
/data/shop/orders.csv(order_id, customer_id, product_id, qty, order_ts, status, channel)、/data/shop/products.csv(product_id, category, price)。売上は、決済完了(status='paid')の注文のqty × priceです。 - スクリプトを
common.pyと同じフォルダー(/root/spk/sql)に置くと、from common import loadが使えます(実行したスクリプトのフォルダーがPythonのモジュールパスに入ります)。 - 結果のCSVは小さいので、
coalesce(1)で1つのファイルにまとめると目で見やすくなります(採点ツールは、ファイルが複数でも読み込みます)。 explainは標準出力に出力します。contextlib.redirect_stdoutで受け取って、ファイルに書き込んでください。- よくあるミス:
rankを使って同点に同じ順位が2つ出ること、累積合計のウィンドウをデフォルト(並び順があるとRANGE … CURRENT ROW)のままにして同じ日が重なること、キャンセルや返金の注文を売上に入れること。 - 公式ドキュメント: Spark SQL Guide・Window Functions・EXPLAIN・Built-in Functions
一時ビューとSQLで月別売上を出す
/root/spk/sql/common.pyに、注文と商品をスキーマを指定して読み込み、一時ビューorders・productsとして登録するload(spark)を作成してください。/root/spk/sql/sql_monthly.pyを、アプリ名spk-sql-monthlyで作成し、spark.sqlで決済完了の注文の月別売上(カラムmonthはyyyy-MM、revenueはqty×priceの合計)を、/root/spk/sql/out/monthly_sqlにヘッダー行ありのCSVで書き込んでください。
一時ビューは、このSparkSessionの中でのみ見える名前です。ビューを登録しても何も読み込まず、プランに名前を付けるだけです。date_format(order_ts, 'yyyy-MM')で月を作り、商品とテーブルを結合して価格を掛けてください。
同じ質問をDataFrame APIで解く
/root/spk/sql/df_monthly.pyを、アプリ名spk-sql-dfで作成し、spark.sqlを使わずにDataFrame API(where・join・groupBy・agg)で、ステップ1と同じ結果を、/root/spk/sql/out/monthly_dfにヘッダー行ありのCSV(month・revenue)で書き込んでください。
F.date_format("order_ts", "yyyy-MM").alias("month")でグループ化し、F.sum(F.col("qty") * F.col("price"))で合計します。結合キーの名前が同じなら、join(p, "product_id")のように文字列で指定して、カラムが1つだけ残るようにしてください。
2つの方法の物理プランが同じか確かめる
/root/spk/sql/plans.pyを、アプリ名spk-sql-plansで作成し、ステップ1・2の2つのクエリをそれぞれ作って、explain(mode="formatted")の出力を、/root/spk/sql/out/plan_sql.txt・/root/spk/sql/out/plan_df.txtに保存したあと、2つのDataFrameをそれぞれcollect()で実際に実行してください。
formattedモードは、上に演算子のツリー、下に番号付きの演算子の説明を出力します。カラム番号(#12のようなもの)はクエリごとに変わりますが、演算子の名前と順序は同じはずです。実際に実行すると、イベントログのSQL実行記録にまったく同じプランのドキュメントが残ります。採点ツールは、自分のファイルがその記録と同じか、2つのプランの演算子が同じかを確認します。
カテゴリ別の上位3商品(ウィンドウによる順位付け)
/root/spk/sql/top3.pyを、アプリ名spk-sql-top3で作成し、カテゴリ(category)・商品別の決済完了の売上を求め、カテゴリごとに売上の降順(同点ならproduct_idの昇順)で1–3位を、/root/spk/sql/out/top3にヘッダー行ありのCSV(category・product_id・revenue・rank)で書き込んでください。
ウィンドウはWindow.partitionBy("category").orderBy(...)です。rank()は同点に同じ番号を付けて次の番号を飛ばすので、ちょうど3個が欲しければ、row_number()に同点のルールを並び順として入れる必要があります。
チャネル別の累積合計(ウィンドウの範囲を決める)
/root/spk/sql/running.pyを、アプリ名spk-sql-runningで作成し、2026年1月の決済完了の注文のチャネル別の日次売上(dayは日付)と、チャネル内で日付順に足した累積合計runningを、/root/spk/sql/out/runningにヘッダー行ありのCSV(channel・day・revenue・running)で書き込んでください。
累積合計は、Window.partitionBy("channel").orderBy("day").rowsBetween(Window.unboundedPreceding, Window.currentRow)上のsumです。日付ごとに1行にまとめてからウィンドウを適用してください。まとめる前に適用すると、同じ日の注文ごとに累積合計が変わってしまいます。
前月比(lagで隣の行を見る)
/root/spk/sql/mom.pyを、アプリ名spk-sql-momで作成し、チャネル別の月売上と、lagで取得した前月売上prev_revenue、増加率growth_pct(=(今月-前月)/前月×100、小数第2位を四捨五入)を、/root/spk/sql/out/momにヘッダー行ありのCSV(channel・month・revenue・prev_revenue・growth_pct)で書き込んでください。チャネルの最初の月は、前月と増加率が空である必要があります。
F.lag("revenue").over(Window.partitionBy("channel").orderBy("month"))が、同じチャネルの直前の月の値を取得します。最初の月は前の行がないためnullで、nullが混ざった割り算もnullになるので、別に処理する必要はありません。
正確なユニーク数と近似ユニーク数
/root/spk/sql/distinct.pyを、アプリ名spk-sql-distinctで作成し、決済が完了した注文をした顧客の数を、countDistinctとapprox_count_distinct(rsd=0.05)で一度に数えて、/root/spk/sql/out/distinct.jsonに{"exact": 정수, "approx": 정수}の形式で書き込んでください(プレースホルダーは順に整数、整数です)。
近似のほうは、HyperLogLog++のスケッチをマージします。ユニークな値をシャッフルで集めなくてよいので、大きくなるほど安くなります。rsdは相対標準誤差の目標値なので、結果が正確な値と何パーセント違うかを自分で計算してみてください。採点ツールは、イベントログのプランに近似関数が実際に入っているかも確認します。
同じプラン・ウィンドウの枠・近似の誤差を残す
/root/spk/sql/report.mdに、## SQL 과 DataFrame・## 윈도 함수・## 근사 집계の3つの節を書いてください(見出しは韓国語で、順に「SQLとDataFrame」「ウィンドウ関数」「近似集計」を意味します)。3つ目の節には、ステップ7の正確な値と近似値を、数値で入れてください。
最初の節には2つのプランで同じだった演算子の名前をいくつか書き、2つ目の節には順位・累積合計・前月比でウィンドウをどう設定したかを1行ずつ書いてください。3つ目の節には、2つの数値とその差のパーセントを書けばかまいません。