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

データパイプライン

深夜3時に落ちた精算ジョブ — 台帳と原子的な差し替え

TT Labで続きを見る

目標

途中で落ちても安全な実行ツールrunner.pyを作ります。出力物を一時的な名前で書いて付け替え、途中まで書かれたファイルが残らないようにし、実行1回を元帳に1行で残し、落ちたシャードから再開し、同じ実行を2回コミットしても数字が増えないようにします。

なぜ重要なのか

パイプラインは必ず途中で落ちます。問題は落ちるという事実ではなく、落ちたときに何が残るかです。宛先ファイルに直接書いていたなら、途中まで書かれたファイルが残り、そのファイルは、サイズも名前も無事なので、次のステップがそのまま読んでいきます。もう一度回すことも危険です。前の実行がどこまで行ったかの記録がなければ、最初からもう一度回るしかなく、最後に帳簿へ追記する場所で、同じ金額が2回足されます。これが「もう一度回したのに、なぜ2倍になったのか」の正体です。このラボで作る仕組みは3つです。1つ目は、アトミックな置換です。同じディレクトリに一時的な名前で書き終えたあと、os.replaceで付け替えます。2つ目は、実行元帳です。実行ごとに1行を追記し、失敗したなら、どのシャードで落ちたかを書きます。3つ目は、再開です。終わったシャードの出力物そのものを目印にして、スキップします。このコースのdp-idempotentのラボは、データベース側の冪等性を扱います。同じ行を2回入れても、1行になるようにすることです。ここは、その前の段階です。プロセスが落ちた場所で、ファイルシステムに何が残るか、残ったものを見て、どこからやり直すかを扱います。採点ツールは、提出された文言を信じません。一時ディレクトリに、採点ツールが作った入力シャードを用意して、自分で作った実行ツールを実際に動かします。わざと落としたあとで、宛先ファイルがそのままか、一時ファイルが宛先の隣に残ったか、元帳に失敗したシャードの名前が書かれたかまで見ます。シャードの数と金額は、実行ごとに変わります。

ステップ

  1. /root/runx/gen_shards.pyを作成して実行し、/root/runx/work/inの下に、入力シャードを作ってください。
  2. /root/runx/runner.pyにscanを作り、シャードの一覧と件数と合計を出力させてください。
  3. partを追加して、シャード1つを処理し、出力物を一時的な名前で書いてから付け替えるようにしてください。--crash=writeで、付け替えの直前に落ちる経路も作ります。
  4. runを追加して、すべてのシャードを処理し、統合した出力物を作らせてください。
  5. 実行ごとに元帳に1行を残し、ledgerで要約を出力させてください。
  6. --crash-shardで、途中のシャードで落ちるようにし、元帳に失敗したシャードが残り、統合した出力物は触られていないことを確認してください。
  7. --resumeを追加して、終わったシャードをスキップし、落ちた場所から続けるようにしてください。
  8. commitを追加して、1日分の帳簿/root/runx/work/out/daily.jsonlに、同じ実行が2回付かないようにしてください。

参考

上流が置いていったシャードを作る

/root/runx/gen_shards.pyを作成して実行し、/root/runx/work/inの下に、シャードファイルを作ってください。シャードは4つ以上、1シャードにつき5行以上、全体で40行以上で、1行はidと整数のamountを含む、JSONの1つの塊です。

シャード1つが、JSON Linesファイル1つです。ファイル名から.jsonlを除いたものがシャード名になるため、ソートしたときに時間順になるように、h00・h01のように桁を合わせて付けてください。乱数のシードを固定しておかなければ、あとで再開をテストするときに、入力が揺れます。

何が入ってきたかを先に数える

/root/runx/runner.pyにscan <작업폴더>(プレースホルダーは作業フォルダです)を作り、シャード名の一覧と、全体の件数と金額の合計を、JSONで出力させてください。シャード名は昇順です。

作業フォルダの下のin/から、.jsonlで終わるファイルだけを選び、名前から拡張子を除きます。空行は数えません。作業フォルダがなければ、終了コード3で終わる必要があり、そうすると、後のステップのエラーメッセージが正直になります。

宛先に直接書かない

part <작업폴더> <조각> [--crash=write](プレースホルダーは、作業フォルダとシャードです)を追加してください。シャードを数えて、out/part-<조각>.json(プレースホルダーはシャード名です)にshard・events・amountを残しますが、宛先と同じディレクトリに一時的な名前で書き終えてから、os.replaceで付け替えます。--crash=writeを渡すと、付け替えの直前に、終了コード9で落ちます。

一時的な名前がpart-*.jsonの一覧に引っかかると、後でそのファイルまで合計に入ります。ドットで始まる名前を使ってください。そして、一時ファイルを/tmpに作ってはいけません。os.replaceはファイルシステムをまたぐと失敗し、その失敗は開発機では再現されません。採点ツールは、落としたあとで、宛先ファイルがそのままか、一時ファイルが宛先の隣に残ったかを見ます。

実行1つにまとめる

run <작업폴더> --run-id=<이름>(プレースホルダーは、作業フォルダと名前です)を追加して、すべてのシャードを順番に処理し、シャードの出力物を読み直して、out/total.jsonにshards・events・amountを残すようにしてください。total.jsonも付け替えで書きます。

総計を、今回処理したシャードだけを足して出すと、あとで再開するときに、スキップしたシャードが抜けます。総計は、常にout/part-*.json全体を読み直して出してください。このルール1つが、再開をただ同然にします。

実行1回を1行で残す

runが終わるときに、/root/runx/work/ledger.jsonlに1行を追記させ、ledger <작업폴더>(プレースホルダーは作業フォルダです)で、{"runs": 정수, "ok": 정수, "failed": 정수, "last": 마지막 줄}(プレースホルダーは、整数、最後の行です)を出力させてください。

元帳は、追記だけをします。前の行を直し始めると、実行1つが1行というルールが崩れ、その瞬間に、元帳はログと変わらなくなります。1行に、run_id・started_at・ended_at・status・shards_total・shards_done・events・amountを入れてください。

途中のシャードで落としてみる

runに--crash-shard=<조각>(プレースホルダーはシャードです)を追加してください。そのシャードの番が来たら、元帳に、statusがfailedでfailed_shardがそのシャードの行を残し、終了コード9で落ちます。統合した出力物out/total.jsonには触れません。

落ちる前に、元帳を先に書く必要があります。元帳がなければ、次の人にできることは、最初からもう一度回すことだけです。shards_doneには、今回の実行で実際に終えたシャードだけを入れ、total.jsonは触らないままにしておきます。前回の実行の答えが、そのまま残っている必要があります。

落ちた場所から続ける

runに--resumeを追加してください。シャードの出力物out/part-<조각>.json(プレースホルダーはシャード名です)がすでにあり、その中のshard名が合っていれば、そのシャードをスキップし、スキップしたものはレスポンスのskippedに、今回処理したものはdoneに入れます。

目印のファイルを別に置かず、シャードの出力物そのものを目印にしてください。出力物は付け替えで作られるため、あるということは、そのシャードが確実に終わったという意味です。目印と出力物を別に置くと、目印だけが残り、出力物は途中まで書かれた状態ができます。総計は、やはりシャードの出力物全体から集め直します。

2回付かないようにする

commit <작업폴더> --run-id=<이름>(プレースホルダーは、作業フォルダと名前です)を追加してください。out/total.jsonを読んで、1日分の帳簿out/daily.jsonlに、run_id・events・amountの1行を付けますが、そのrun_idがすでにあれば付けず、{"appended": false, ...}を出力します。自分の作業フォルダでも、実行を1つコミットしておいてください。

シャードの処理は上書きなので、何回やっても同じですが、帳簿に1行付けることは、呼ぶたびに増えます。重複を時刻で判断すると、同じ日に2回回った実行を区別できません。実行の名前で判断してください。