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

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

パーティション化したレイクを書き、上書きし、小さなファイルを減らす

TT Labで続きを見る

目標

注文を日付で分けたParquetのレイクとして書き、4つの保存モード(errorifexists・append・overwrite・ignore)のうち3つが実際に何をするかを確認します。静的上書きと動的パーティション上書きの違いをディレクトリ数で見て、小さなファイルを減らす2つのノブ(パーティションカラムでまとめる・ファイルあたりの行数の上限)を使ってみます。

なぜ重要なのか

Sparkの書き込みは、一度で終わる処理ではありません。タスクごとに一時パスに書き、ジョブが成功するとコミットプロトコルが結果を元の場所に移して、_SUCCESSマーカーを残します。そのため、書きかけの結果は読まれません。その代わり、保存モードを誤ると、この正直な仕組みが、正直に事故を起こします。 appendは、再試行1回で同じデータを2回入れます。overwriteは、デフォルトが静的なので、1日分だけ直そうとして書いたものが、対象パスのすべての日付を消します。動的パーティション上書き(spark.sql.sources.partitionOverwriteMode=dynamic)は、書き込むデータに含まれるパーティションだけを入れ替えます。この違いを知らずに本番のレイクにoverwriteを1回使うと、数か月分が一度に消えます。 ファイル数は、書き込むタスク数×そのタスクが持つパーティション値の数です。シャッフル後のタスクがすべての日付を少しずつ持っていると、日付ディレクトリごとに、タスク数だけ小さなファイルができます。パーティションカラムで先にまとめたり(repartition("칼럼")、プレースホルダーはカラム名です)、ファイルあたりの行数に上限(maxRecordsPerFile)を設けたりすれば、サイズを制御できます。

ステップ

  1. /root/spk/write/common.pyに、注文を読み込んでorder_dateを付ける関数を置き、/root/spk/write/lake.py(アプリspk-write-lake)で、/root/spk/write/lake/ordersにorder_dateで分けたParquetを書き込んでください。
  2. /root/spk/write/exists.py(アプリspk-write-exists)で、同じパスにモードを指定せずにもう一度書き込ませ、エラー条件名を、/root/spk/write/out/exists.txtに書き込んでください。
  3. /root/spk/write/append.py(アプリspk-write-append)で、1月の注文を、/root/spk/write/lake/appendにoverwriteで1回、appendでもう1回書き込み、読み直した行数とユニークなorder_idの数を、/root/spk/write/out/append.jsonに書き込んでください。
  4. /root/spk/write/static.py(アプリspk-write-static)で、全体を、/root/spk/write/lake/staticに書き込んだあと、2026-02-10の1日分だけをoverwriteでもう一度書き込み、残ったパーティションディレクトリの数を、/root/spk/write/out/static.jsonに書き込んでください。
  5. /root/spk/write/dynamic.py(アプリspk-write-dynamic、partitionOverwriteMode=dynamic)で、全体を、/root/spk/write/lake/dynamicに書き込んだあと、2026-02-10からキャンセル注文を除いたものだけをoverwriteで書き込み、残ったパーティションディレクトリの数を、/root/spk/write/out/dynamic.jsonに書き込んでください。
  6. /root/spk/write/small.py(アプリspk-write-small)で、クリックをrepartition(40)して、/root/spk/write/lake/clicks_manyに、repartition("page")して、/root/spk/write/lake/clicks_fewに、それぞれpageで分けて書き込み、2つのレイクのデータファイル数を、/root/spk/write/out/files.jsonに書き込んでください。
  7. /root/spk/write/maxrec.py(アプリspk-write-maxrec)で、注文をrepartition("channel")したあと、maxRecordsPerFile=20000を指定して、/root/spk/write/lake/by_channelにchannelで分けて書き込んでください。
  8. /root/spk/write/report.mdに、## 저장 모드・## 정적과 동적 덮어쓰기・## 파일 수の3つの節を書いてください(見出しは韓国語で、順に「保存モード」「静的上書きと動的上書き」「ファイル数」を意味します)。2つ目の節にはステップ4・5のパーティション数を、3つ目の節にはステップ6の2つのファイル数を入れてください。

参考

日付で分けたレイクを書く

/root/spk/write/common.pyに、注文をスキーマを指定して読み込み、order_date = to_date(order_ts)を付ける関数を置き、/root/spk/write/lake.pyを、アプリ名spk-write-lakeで作成して、/root/spk/write/lake/ordersにpartitionBy("order_date")でParquetを書き込んでください(mode("overwrite"))。

日付90日なら、ディレクトリは90個です。ジョブが成功して初めて_SUCCESSができます。このマーカーがないディレクトリは書きかけの可能性があるので、読む側がまず確認します。採点ツールは、ディレクトリ数・行数・マーカーを確認します。

デフォルトのモードは止まる(errorifexists)

/root/spk/write/exists.pyを、アプリ名spk-write-existsで作成し、ステップ1と同じパスにモードを指定せずにもう一度書き込ませて、例外のgetCondition()を、/root/spk/write/out/exists.txtの1行目に書き込んでください。ステップ1のレイクはそのまま残っている必要があります。

デフォルトの保存モードは「すでにあればエラー」です。誤って同じパスに書き込むのを防ぐ安全装置であり、本番のジョブでモードを明示する習慣が必要な理由でもあります。

append(再試行が2倍を作る)

/root/spk/write/append.pyを、アプリ名spk-write-appendで作成し、2026年1月の注文を、/root/spk/write/lake/appendにoverwriteで1回、appendでもう1回書き込み、読み直した行数とユニークなorder_idの数を、/root/spk/write/out/append.jsonに{"rows": 정수, "distinct_orders": 정수}の形式で書き込んでください(プレースホルダーは順に整数、整数です)。

appendは、すでにあるファイルを見ずに新しいファイルを追加します。同じ入力でジョブを再試行すると、行が2倍になり、ファイル名が違うので、見た目には問題なく見えます。冪等に書き込むには、上書き(できれば動的)か、キーでマージする方式を使います。

静的上書き(1日を直そうとして全部を消す)

/root/spk/write/static.pyを、アプリ名spk-write-staticで作成し、全注文を、/root/spk/write/lake/staticにorder_dateで分けて書き込んだあと、2026-02-10の1日分だけを、同じパスにoverwrite(設定はデフォルト)でもう一度書き込んでください。残ったorder_date=ディレクトリの数を、/root/spk/write/out/static.jsonに{"partitions_after": 정수}の形式で書き込んでください(プレースホルダーは整数です)。

partitionOverwriteModeのデフォルトはstaticです。overwriteは書き込みの前に、対象パスの下をまるごと消します。書き込むデータにどんな日付が含まれているかは見ません。このステップの結果は事故です。それを目で見ることが目的です。

動的パーティション上書き(その日だけを変える)

/root/spk/write/dynamic.pyを、アプリ名spk-write-dynamic、設定spark.sql.sources.partitionOverwriteMode=dynamicで作成し、全注文を、/root/spk/write/lake/dynamicにorder_dateで分けて書き込んだあと、2026-02-10の注文からキャンセル(cancelled)を除いたものだけを、同じパスにoverwriteで書き込んでください。残ったorder_date=ディレクトリの数を、/root/spk/write/out/dynamic.jsonに{"partitions_after": 정수}の形式で書き込んでください(プレースホルダーは整数です)。

動的モードでは、書き込むデータに含まれるパーティション(ここでは1日)だけを差し替えます。残りの89日はそのままです。採点ツールは、ディレクトリ数、その日にキャンセル注文がないか、ほかの日が元のままかを確認します。

小さなファイル(タスク数×パーティション値)

/root/spk/write/small.pyを、アプリ名spk-write-smallで作成し、/data/clicks/clicks.jsonlをrepartition(40)して、/root/spk/write/lake/clicks_manyに、repartition("page")して、/root/spk/write/lake/clicks_fewに、それぞれpartitionBy("page")で書き込み、2つのレイクのデータファイル(page=*/part-*)の数を、/root/spk/write/out/files.jsonに{"many": 정수, "few": 정수}の形式で書き込んでください(プレースホルダーは順に整数、整数です)。

書き込みタスク1つは、自分が持つパーティション値ごとにファイルを1つずつ開きます。タスク40個がページ7つをすべて持っていれば、最大280個です。パーティションカラムで先にまとめると、ページ1つがタスク1つに入り、ページあたりファイル1つになります(その代わり、1ページが大きいと、そのタスクが重くなります)。

ファイルあたりの行数に上限を設ける

/root/spk/write/maxrec.pyを、アプリ名spk-write-maxrecで作成し、注文をrepartition("channel")したあと、option("maxRecordsPerFile", 20000)で、/root/spk/write/lake/by_channelにpartitionBy("channel")で書き込んでください。

チャネルでまとめると、チャネル1つがタスク1つに入ってファイル1つになりますが、最も大きいチャネルは17万行を超えます。上限を設けると、タスクが2万行ごとにファイルを新しく開きます。採点ツールは、チャネルごとのファイル数が切り上げ(行数÷20000)か、ファイルごとに行が2万を超えないかを、Parquetメタデータで確認します。

書き込みのルールをチームのルールにする

/root/spk/write/report.mdに、## 저장 모드・## 정적과 동적 덮어쓰기・## 파일 수の3つの節を書いてください(見出しは韓国語で、順に「保存モード」「静的上書きと動的上書き」「ファイル数」を意味します)。2つ目の節にはステップ4・5の残りのパーティション数を、3つ目の節にはステップ6の2つのファイル数を入れてください。

本番のレイクに書き込むジョブなら、どのモードをデフォルトにして何を禁止するかを、ルールとして書いてみてください。数値は、自分の結果ファイルから書き写してください。