Apache Spark — 遅いジョブの答えは実行計画とイベントログにある
変換は計画を積むだけで、アクションがジョブを起動してステージとタスクに分ける
一言でいうと
Sparkのコードは上から下へ順に実行されているように見えますが、実際はドライバーがプランを積み上げ、アクション(action)に出会った瞬間にそのプラン全体をジョブ(job)に変換し、ステージとタスクに分割してエグゼキューターに割り振ります。何がいつ起きたかを教えてくれるのは、コードではなくイベントログです。
なぜこう分かれているのか
Sparkのコードを初めて読むと、spark.read.csv(...)の行でファイルを読み、filterの行で絞り込み、countの行で数える、と考えがちです。そう読むと、性能問題を見当違いの場所で探すことになります。「読み込みに40秒かかる」ように見える箇所は、実際にはそれまでに積み上げた20個のトランスフォーメーションが一度に動く箇所です。
Sparkがこのように仕事を分けているのは、データが1台のメモリに収まらないという前提から出発したからです。計算を複数のプロセスに分担させるには、誰かが全体像を把握して仕事を細かく切り分ける必要があり、受け取った側は自分の担当分だけを計算すればよいのです。クラスターモード概要の用語表は、この役割を次のように定義しています。
- ドライバー(driver): アプリのmain()を実行し、SparkContextを作成するプロセスです。プランを立て、タスクを割り振ります。
- エグゼキューター(executor): ワーカーノードにアプリごとに起動するプロセスで、タスクを実行し、データをメモリやディスクに保持します。
- タスク(task): エグゼキューター1つに送る仕事の単位です。
- ジョブ(job): アクション1つに応答して生じる、複数のタスクからなる並列計算です。
- ステージ(stage): ジョブを、互いに依存するより小さなタスクの集まりに分けたものです。ドキュメントは、MapReduceのマップ・リデュースの段階に似ていると書いています。
ジョブの定義に「アクションに応答して」という言葉が入っていることが核心です。アクションがなければジョブもありません。
どう動くのか
RDDプログラミングガイドは、「Sparkのすべてのトランスフォーメーションは遅延評価される」と明記しています。トランスフォーメーションは結果をすぐには計算せず、どのトランスフォーメーションを適用したかだけを記憶し、アクションがドライバーに結果を返すよう求めたときに初めて計算されます。ドキュメントは、この設計の利点も併せて書いています。mapのあとにreduceが来ると事前にわかっていれば、大きな中間結果の代わりに、縮約した結果だけをドライバーに返せます。最後まで見てから道筋を決められることが、遅延評価の価値です。
df = spark.read.csv("/data/shop/orders.csv", header=True, schema=ddl) # 계획 1줄
paid = df.filter(F.col("status") == "paid") # 계획 +1
big = paid.withColumn("bulk", F.col("qty") >= 5) # 계획 +1
big.count() # 행동 — 여기서 잡이 뜨고, 위 세 줄이 한 번에 실행된다
アクションに出会うと、ドライバーはプランを物理プランに変換し、それをシャッフル境界で切り分けます。シャッフルのない区間(読み込み・絞り込み・カラム追加のようなナロートランスフォーメーション)は、1つのタスクの中で連続して実行できるので1つのステージになります。groupByのように同じキーを1か所に集める必要があるワイドトランスフォーメーションが出てくると、前のステージのすべてのタスクが結果をシャッフルファイルに書き終えてからでないと後ろのステージがそれを読めないため、そこでステージが分かれます。
ステージの中ではパーティション1つがタスク1つです。つまり、パーティション数がそのまま同時に作業できる断片の数になります。ファイルを読むときのパーティション数は、ファイルサイズだけでは決まりません。パフォーマンスチューニングのドキュメントを見ると、ファイルを分割する最小パーティション数(spark.sql.files.minPartitionNum)のデフォルト値はspark.sql.leafNodeDefaultParallelismで、その値のデフォルトはSparkContextのデフォルト並列度です。小さなファイルでも、コア数の分だけは分けて読もうとするという意味です。
localモードではどれがドライバーでどれがエグゼキューターか
このラボのPodにはクラスターがありません。マスターURLの表では、localはワーカースレッド1つ(並列なし)、local[K]はワーカースレッドK個です。localモードではドライバーのJVM 1つがエグゼキューターの役割も担います。タスクはそのJVM内のスレッドとして動きます。
それでも、ドライバー・エグゼキューター・ジョブ・ステージ・タスクの区別はそのまま残ります。プランを立てるコードとタスクを実行するコードが、同じプロセスにあるだけです。localモードで学ぶことが無駄にならない理由がここにあります。イベントログに残るジョブ・ステージ・タスクのイベントの形は、クラスターの場合と同じです。このPodはlocal[2]、ドライバーメモリ1gに設定してあります。
終了したアプリを振り返る証拠としてのイベントログ
ドライバーのWeb UI(4040)は、アプリが動いている間しか見られません。終了したアプリを見るには、モニタリングのドキュメントのとおりヒストリーサーバーを使いますが、ヒストリーサーバーはアプリが残したイベントログでUIを再構築します。イベントログがなければ、終了したアプリについてわかることはほとんどありません。
イベントログは、JSON 1行が1つのイベントです。ジョブが始まるとSparkListenerJobStart、ステージが終わるとSparkListenerStageCompletedが記録されます。そのため、「アクション3回ならジョブも3つ」のような推測を行数を数えて検証できます。
注意点が1つあります。設定のドキュメントによると、Spark 4.2のデフォルトは、spark.eventLog.compressがtrue、圧縮コーデックがzstd、spark.eventLog.rolling.enabledがtrueです。デフォルトのままにすると、ログは複数のファイルにローリングされて圧縮されるため、grepですぐには読めません。このラボでは、プレーンテキストの1ファイルに変更してあります。本番で同じ方法を使うなら、この2つの設定をまず確認する必要があります。
現場での姿
アクション1つがジョブ1つとは限りません。スキーマを指定せずにヘッダー行付きのCSVを読むと、Sparkはヘッダー行を確認するために小さなジョブを先に起動します。ファイルを書き込むアクションも、ジョブを1つより多く作ることがあります。そのため、コードを数えるのではなくログを数えます。このラボでも、アクションを3回呼んだアプリのジョブ数を推測せず、ログから直接数えます。
エラーが出る場所は2通りあります。パスが間違っていると、アクションを待たずにreadした瞬間、分析の段階で失敗します。このときのエラー条件名がPATH_NOT_FOUNDです。一方、データの中の値が問題なら、プランを立てる時点ではわからず、アクションが実際にデータを読むときになって初めて失敗します。前者はコードの行とエラーが一致し、後者は見当違いの行(アクション)で発生します。
ドライバーが先に落ちます。結果を目で見ようとcollect()を呼ぶと、結果全体がドライバー1か所に集まります。RDDガイドは、これがドライバーのメモリを溢れさせる可能性があるため、数件だけ見たいときはtakeを使うよう勧めています。エグゼキューターがいくら多くても、ドライバーは1つです。
トランスフォーメーション内のprintはドライバーの画面に出ません。トランスフォーメーション内のコードはエグゼキューターで動くからです。localモードでは偶然見えますが、クラスターに移すと消えます。
実務で本当に大切なこと
- トランスフォーメーションはプランであり、アクションが仕事を実行させます。遅い箇所は多くの場合アクションのある行ですが、原因はその前のトランスフォーメーション全体です。
- ステージの境界はシャッフルです。ステージ数を減らすことは、そのままシャッフルを減らすことです。
- パーティション1つがタスク1つです。パーティション数が、同時に作業できる上限を決めます。
- ジョブ数は推測せず、イベントログから数えます。スキーマ推論・ヘッダー行の確認・書き込みが、ジョブを増やします。
- イベントログのデフォルトはローリングとzstd圧縮です。読むツールを選ぶ前に、まず設定を確認します。
- ドライバーは1つです。collectで結果をドライバーに集めることは、最もよくあるドライバーのメモリ事故です。
次のラボですること
spark-submitでバージョンを確認したあと、注文30万件を数える最初のジョブを起動します。トランスフォーメーションだけを積んだアプリを実行してイベントログにジョブが1つもないことを確認し、アクションを複数回呼んだアプリでジョブが実際にいくつ起動したかをログから数えます。groupByでシャッフルを起こしてステージが分かれる様子を見て、local[1]とlocal[2]で同じファイルのパーティション数がどう変わるかを比べます。最後に、存在しないパスを読んでエラー条件名を受け取り、数値をレポートにまとめます。