Apache Spark — 遅いジョブの答えは実行計画とイベントログにある
最初のジョブを起動してイベントログで確かめる — 変換は遅延評価、アクションが処理を起こす
目標
localモードのSparkで注文30万件を読み込んで数え、トランスフォーメーションだけを積んだアプリとアクションを呼んだアプリをイベントログで比べます。ジョブ・ステージ・タスクがいつ生じ、パーティション数が何で決まるのかを、数値で確認します。
なぜ重要なのか
Sparkのコードを初めて読むと、1行1行がその場で実行されるように見えます。そうではありません。filter・withColumn・selectのようなトランスフォーメーションはプランに1行を加えるだけで、count・take・writeのようなアクションを呼んだときに、そのプラン全体がジョブに変換されて実行されます。これを知らないと、「読み込みが遅い」ように見える箇所が、実はそれまでに積み上げたトランスフォーメーション全体が一度に動く箇所だという点を見落とします。
その違いは目で確認できます。Sparkはアプリが終了すると、何をしたかをイベントログとして残します。ジョブがいくつ起動したか、ステージがいくつあったか、タスクがいくつあったか、どんな物理プランを使ったかが、すべてJSON 1行ずつに入っています。本番で終了したジョブを振り返るときに見るのもこれです(ヒストリーサーバーがこのファイルを読みます)。
localモードでは、ドライバーのJVM 1つがそのままエグゼキューターです。local[2]はコアを2つ使うという意味で、この数字がデフォルトの並列度と、ファイルを何個の断片に分けるかを同時に決めます。
ステップ
- /root/spk/first/version.txtに、
spark-submit --versionの出力(標準エラー出力を含む)を保存してください。 - /root/spk/first/count.pyを作成し、アプリ名
spk-first-countで/data/shop/orders.csv(ヘッダー行あり)の行数を数えて、/root/spk/first/out/count.jsonに{"rows": 정수}の形式で書き込んでください(プレースホルダーは整数です)。 - /root/spk/first/lazy.pyを作成し、アプリ名
spk-first-lazyでスキーマを直接指定して読み込み、filter(またはwhere)とwithColumnを積み重ねてください。ただしアクションは呼ばないでください。実行後、そのアプリのジョブが0個である必要があります。 - /root/spk/first/actions.pyを作成し、アプリ名
spk-first-actionsでアクションを3回以上(例:count・take・write)呼び出し、/root/spk/first/out/sampleに100行をJSONで書き込んでください。そのアプリのイベントログでジョブ開始イベントの数を数え、/root/spk/first/out/jobs.txtに整数1つで記入してください。 - /root/spk/first/agg.pyを作成し、アプリ名
spk-first-aggでstatus別の注文数を、/root/spk/first/out/by_statusにヘッダー行ありのCSVで書き込んでください。そのアプリで完了したステージ数を、/root/spk/first/out/stages.txtに整数で記入してください。 - /root/spk/first/par.pyを作成して最初の引数をアプリ名として受け取るようにし、
--master local[1]でspk-first-par1を、--master local[2]でspk-first-par2を実行して、結果をそれぞれ、/root/spk/first/out/spk-first-par1.json・/root/spk/first/out/spk-first-par2.jsonに、master・default_parallelism・partitions(注文ファイルを読み込んだDataFrameのパーティション数)として書き込んでください。 - /root/spk/first/fail.pyを作成し、アプリ名
spk-first-failで存在しないパス/data/shop/order.csvを読み込ませてください。捕捉した例外のエラー条件名(getCondition())を、/root/spk/first/out/error.txtの1行目に書き込んでください。 - /root/spk/first/report.mdに、
## 잡과 스테이지・## 게으른 실행・## 파티션の3つの節で、自分の数値を書いてください(見出しは韓国語で、順に「ジョブとステージ」「遅延実行」「パーティション」を意味します)。最初の節にはステップ4のジョブ数を、3つ目の節にはステップ6の2つのパーティション数を、数値で入れてください。
参考
- 実行は
spark-submit /root/spk/first/count.pyのようにします。基本設定は/opt/spark/conf/spark-defaults.confにあります(local[2]、ドライバーメモリ1g、イベントログ有効)。 - イベントログは、アプリが終了すると
/root/spark-events/local-<숫자>として残ります(プレースホルダーは数字です)。実行中は、名前の末尾に.inprogressが付きます。どのファイルがどのアプリのものかは、grep -l '"App Name":"spk-first-actions"' /root/spark-events/*で探します。 - ジョブ数は
grep -c '"Event":"SparkListenerJobStart"' <파일>(プレースホルダーはファイルです)で、ステージ数は"Event":"SparkListenerStageCompleted"を数えます。jq -c 'select(.Event=="SparkListenerJobStart")'で中身を見ることもできます。 - よくあるミス:
spark.stop()を抜かしてログが.inprogressのまま残ること、ヘッダー行付きのCSVでheader=Trueを抜かしてヘッダー行を1行として数えること、スキーマなしで読んでヘッダー行を確認するジョブが生じること(ステップ3はジョブが0個でなければなりません)。 - 公式ドキュメント: Cluster Mode Overview・RDD Programming Guide — RDD Operations・Monitoring — Viewing After the Fact・Submitting Applications — Master URLs
Sparkのバージョンを確認する
spark-submit --versionの出力を標準エラー出力も含めて、/root/spk/first/version.txtに保存してください。
Sparkはバージョン情報を標準エラー出力に出します。2>&1で合流させないと、ファイルが空になってしまいます。バージョン番号とともに、どのScalaやJava上で動いているかも表示されます。Spark 4はJava 17以上を要求します。
最初のジョブ(30万件を数える)
/root/spk/first/count.pyを、アプリ名spk-first-countで作成し、/data/shop/orders.csv(ヘッダー行あり)の行数を数えて、/root/spk/first/out/count.jsonに{"rows": 정수}の形式で書き込んでください(プレースホルダーは整数です)。spark-submitで実行してください。
SparkSession.builder.appName(...)がアプリ名を決めます。採点ツールは、その名前のイベントログが最後まで閉じられているかと、ファイルの数値が元データの行数と一致しているかを確認します。ヘッダー行を行として数えると、1つ多くなります。
トランスフォーメーションだけ積んでもジョブは起動しない
/root/spk/first/lazy.pyを、アプリ名spk-first-lazyで作成し、スキーマを直接指定したspark.read.csvで注文を読み込み、filter(またはwhere)とwithColumnを積み重ねてください。アクションは1つも呼ばないでください。実行後、そのアプリのイベントログにジョブが0個である必要があります。
スキーマを指定しないと、Sparkはヘッダー行を確認するために小さなジョブを起動します。schema="order_id STRING, ..."のようにDDL文字列を渡せば、そのジョブも消えます。スキーマを出力するのは問題ありません。スキーマはドライバーがすでに把握しているので、データを読み込みません。
アクションごとにジョブが起動する(イベントログで数える)
/root/spk/first/actions.pyを、アプリ名spk-first-actionsで作成し、アクションを3回以上呼び出し、そのうち1つで100行を、/root/spk/first/out/sampleにJSONで書き込んでください。実行後、そのアプリのイベントログでSparkListenerJobStartイベントの数を数え、/root/spk/first/out/jobs.txtに整数1つで記入してください。
アクション1つがジョブ1つだと決めつけないでください。ファイルを書き込むアクションやスキーマ推論のように、ジョブを2つ以上作るアクションもあります。そのため、数えるのはコードではなくログです。同じ名前で何度も実行した場合は、最新のログを数える必要があります。
ワイドトランスフォーメーションがステージを分ける
/root/spk/first/agg.pyを、アプリ名spk-first-aggで作成し、status別の注文数を、/root/spk/first/out/by_statusにヘッダー行ありのCSV(カラムはstatusとcount)で書き込んでください。そのアプリのイベントログでSparkListenerStageCompletedイベントの数を数え、/root/spk/first/out/stages.txtに整数で記入してください。
groupByは同じキーを1つのタスクに集める必要があるのでシャッフルが生じ、シャッフルの前後が別々のステージになります。採点ツールは、CSVが元データで数えた値と一致しているか、そのアプリでシャッフルを使ったステージが実際にあったか、記入したステージ数がログと一致しているかを確認します。
コア数がパーティション数を決める
/root/spk/first/par.pyを、最初の引数をアプリ名として受け取るように作成し、spark-submit --master 'local[1]' par.py spk-first-par1とspark-submit --master 'local[2]' par.py spk-first-par2で2回実行して、/root/spk/first/out/spk-first-par1.json・/root/spk/first/out/spk-first-par2.jsonに{"master": 문자열, "default_parallelism": 정수, "partitions": 정수}を書き込んでください(プレースホルダーは順に文字列、整数、整数です)。partitionsは、スキーマを指定して読み込んだ注文DataFrameのrdd.getNumPartitions()です。
ファイルを何個の断片に分けるかは、spark.sql.files.maxPartitionBytes(デフォルト128MB)だけでは決まりません。コアが複数あるとき、全体のサイズをコア数で割った値のほうが小さければ、そのサイズで切ります。コアを遊ばせないためのルールです。16MBのファイルが、コア1つと2つでどう変わるかを見てください。
失敗もアプリである(エラー条件名を読む)
/root/spk/first/fail.pyを、アプリ名spk-first-failで作成し、存在しないパス/data/shop/order.csvを読み込ませ、例外を捕捉してgetCondition()が返したエラー条件名を、/root/spk/first/out/error.txtの1行目に書き込んでください。アプリはspark.stop()で正常に終了させてください。
存在しないパスは、アクションを呼んだときではなく、readした瞬間に明らかになります。ファイル一覧はプランを立てるときに確認するからです。Spark 4のエラーは、人間が読む文章のほかに、機械が読む名前(エラー条件)を持っています。アラートやリトライのルールは、文章ではなくこの名前で設定するのが安全です。
見たものを数値で残す
/root/spk/first/report.mdに、## 잡과 스테이지・## 게으른 실행・## 파티션の3つの節を書いてください(見出しは韓国語で、順に「ジョブとステージ」「遅延実行」「パーティション」を意味します)。最初の節にはステップ4で数えたジョブ数を、3つ目の節にはステップ6の2つのパーティション数を、数値で入れてください。
数値を書き写すだけでなく、なぜその数値になるのかを1、2文ずつ添えてください。アクション3回でジョブがいくつだったか、トランスフォーメーションだけを積んだアプリはなぜ0個だったか、コア数がパーティション数をどう変えたか。これがこのラボで見たすべてです。