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

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

書き込みは取り消せないので、モードとファイルの形を先に決める

TT Labで続きを見る

一言でいうと

Sparkの書き込みは、保存モード(すでにあるときにどうするか)、パーティションディレクトリ(どんな形で置くか)、コミットプロトコル(書き終えたことをどう知らせるか)の3つで決まり、3つとも選び方を誤ると、エラーなしにデータが消えたり増えたりします。

なぜ書き込みを別に学ぶのか

読み込みやトランスフォーメーションは、何度間違えても元データはそのままです。書き込みは違います。上書き1回で3か月分のパーティションが消え、再試行したappendが同じ1日を2回入れます。そして、どちらも成功で終わります。ジョブは緑のままで、レポートだけが間違います。

ロード・保存のドキュメントは、保存モードを説明しながら、この点をまず警告しています。保存モードはロックを使わず、アトミックでもありません。そして、上書きは新しいデータを書く前に既存のデータを消します。上書きの途中でジョブが死ぬと、古いデータも新しいデータもない瞬間が生じるという意味です。

どう動くのか(4つの保存モード)

同じドキュメントが挙げるモードは4つです。

デフォルトがエラーであるのは、良い設計です。何をするか指定していない書き込みが、他人のデータに手を出せないようにするからです。問題は、人々がそのエラーをなくそうとして、習慣のようにoverwriteを付けることから始まります。

パーティションディレクトリとファイル数

partitionBy("day")で書くと、値ごとにday=2026-01-03/のようなディレクトリができ、そのカラムはファイルの中ではなくパスに入ります。読む側は、パスから値を復元し、そのカラムに条件がかかれば、ディレクトリごとスキップします。同じドキュメントは、この方式がディレクトリ構造を作るため、値が非常に多いカラムには向かないと書いています。ユーザーIDでpartitionByすると、ディレクトリがユーザー数だけできます。

ファイル数は、もっと静かな落とし穴です。書き込みタスク1つは、自分が持つ行のパーティション値ごとに1ファイルを書きます。シャッフル後の200個のタスクがすべて30日分を少しずつ持っていれば、ファイルは最大で200 × 30 = 6,000個になります。毎日、数KBのファイルが200個です。

(df.repartition("day")                 # 같은 날은 한 태스크로 — 날마다 파일 하나
   .write.partitionBy("day")
   .option("maxRecordsPerFile", 50000) # 너무 큰 날은 5만 줄씩 끊는다
   .mode("overwrite")
   .option("partitionOverwriteMode", "dynamic")
   .parquet("/data/lake/orders"))

書き込みの前にパーティションカラムでrepartitionすると、同じ日の行が1つのタスクに集まり、毎日ファイル1つになります。その代わり、大きい日が巨大なファイル1つになります。その上限を設けるのが、設定ドキュメントのspark.sql.files.maxRecordsPerFileです。1ファイルに書く最大レコード数で、デフォルトの0は制限なしという意味です。書き込みオプションとして指定することもできます。

値が多いカラムで分けたいときの代替案が、バケットです。同じドキュメントによると、bucketByは値の種類数に関係なく、決まった数のバケットにデータをハッシュで振り分けるので、ユニークな値が際限なく増えるカラムにも使えます。その代わり、バケットとソートは永続テーブル(saveAsTable)にだけ適用されます。パスにファイルを書くだけのsave()では、バケットを残せません。日付のように種類が少なく、検索条件に常に入るカラムはpartitionBy、ユーザーIDのように種類が多く結合キーとして使われるカラムはbucketByと、分けて覚えればよいです。

静的上書きと動的上書き

パーティション化されたパスにoverwriteすると、何が消えるでしょうか。設定ドキュメントのspark.sql.sources.partitionOverwriteModeがこれを決め、デフォルトはSTATICです。静的モードは、書き込みの前に、対象に該当するパーティションをあらかじめ消します。パス全体にDataFrameを上書きすると、対象はパス全体です。1日分だけ直そうとして、その日の行だけを含むDataFrameで上書きしたところ、残りの日付がすべて消える事故が、これです。

動的(dynamic)モードは、あらかじめ消さずに、実際にデータが書かれたパーティションだけを入れ替えます。1日分のDataFrameで上書きすれば、その1日だけが変わり、ほかの日はそのままです。ドキュメントは、書き込みオプションのpartitionOverwriteModeがセッション設定より優先されると書いています。そのため、上の例のように、書き込む場所にオプションとして書いておくほうが安全です。セッション設定はほかの人が変えられますが、コードに書かれたオプションは、その書き込みと一緒に動きます。

コミットプロトコルと_SUCCESS

複数のタスクが1つのディレクトリに同時に書くとき、読む側は、いつ書き終えたものを信頼できるのでしょうか。タスクは最終的な場所ではなく、_temporary/の下の自分の場所に書きます。タスクが成功すると、その出力がコミットされ、すべてのタスクが終わると、ジョブコミットが出力を最終的な場所に移します。Hadoopのmapred-default.xmlは、アルゴリズムバージョン1のジョブコミットが、タスクの出力をまとめて移し、_temporaryを消したあと、_SUCCESSを書くと説明しています。そのため、_SUCCESSは「このディレクトリは最後まで書かれた」というマーカーです。失敗したジョブのディレクトリにはありません。

Sparkは、このコミットアルゴリズムを設定で選びます。設定ドキュメントでspark.hadoop.mapreduce.fileoutputcommitter.algorithm.versionのデフォルトは1で、2はMAPREDUCE-7282のような正確性の問題を起こすことがあると書かれています。バージョン1は、移動(リネーム)に頼ります。クラウド連携のドキュメントは、オブジェクトストレージではリネームが非常に遅く、失敗すると状態がわからなくなると警告しています。ローカルディスクやHDFSでは安価なことが、S3では高価なことになる理由です。

現場での姿

1つ目は、1日を直そうとして全部を消すことです。静的上書きの典型です。パーティション化されたレイクを上書きするコードには、動的モードをオプションとして埋め込んでおきます。

2つ目は、再試行が2倍を作ることです。appendで書くジョブが途中で失敗して再実行されると、先にコミットされた部分の上に、同じデータがまた入ります。1日単位で再実行できるジョブは、appendよりも、その日のパーティションの動的上書きのほうが安全です。何度実行しても結果が同じになるからです。

3つ目は、次のジョブが書きかけのディレクトリを読むことです。後続のジョブが、ファイルが見えたとたんに読み始めると、書き込み途中の状態を見てしまいます。_SUCCESSを確認してから読むというルール1つが、これを防ぎます。パーティションディレクトリ1つ1つではなく、書き込みの最上位のパスに1つできる点も、覚えておきます。読む側は、日付ディレクトリではなく、レイクの最上位でこのマーカーを探します。

実務で本当に大切なこと

次のラボですること

注文を日付でpartitionByしてレイクを書き、同じパスにデフォルトのモードでもう一度書いたときに出るエラー条件を確認します。overwriteのあとappendでもう一度書いて、行が2倍になり、ユニークな注文番号はそのままであることを数え、1日分で静的上書きをしたときにほかの日付が消えるのを見たあと、動的上書きでその1日だけを直します。最後に、クリックを40個の断片に散らして書いたレイクと、パーティションカラムで再分割して書いたレイクのファイル数を比べ、maxRecordsPerFileで、1ファイルに入るレコード数に上限をかけます。