最後に確定した位置から取り込みを再開する:設計原理
一言でいうと
原本のフィンガープリントとチェックポイントを組み合わせて、別のファイルで誤って引き継ぐことを防ぎます。
なぜ必要なのか
数千行をインポートしていた作業が、途中で死にました。運用者がファイルを差し替えて、同じ作業idでもう一度実行したところ、前半は古いファイル、後半は新しいファイルという結果ができました。処理した行番号だけを保存していると、入力の同一性を確認できません。原本のバイト列のフィンガープリントと、最後にコミットした位置を、一緒に保持する必要があります。
どう動くのか
入力は、重複するidがないJSON配列で、各行はidとvalueを持ちます。原本のbytesのSHA-256とチェックポイントを、importsに保存します。バッチの各itemの挿入とnext_indexの増加は、同じトランザクションです。途中の行で例外が起きたら、そのバッチ全体がロールバックされ、以前に終わったバッチは残ります。再実行は、保存されたインデックスから始めますが、原本のフィンガープリントが違えば拒否します。
원본 bytes → 지문 확인 → next_index → 배치 INSERT + 체크포인트 COMMIT
└ 중간 실패 → 이번 배치만 ROLLBACK
契約を読んで失敗を予測するワークシート
以下は、実装を丸ごと暗記するための解答ではなく、ステップごとのコードレビューです。各変更の断片は、意図的に契約を壊しています。変更後も、正常なケースは通ることがある点に注意してください。実行する前に、どの入力・例外・状態を観測すれば違いが現れるかを予想し、実装したあとで、その予想と結果を比べます。
1. 入力行の契約を確認する
parse_rows(raw)は、JSON配列のbytesを読みます。各項目は、id(空でないstr)とvalue(boolを除くint)を持ち、idの重複は禁止します。違反はValueErrorです。{id,value}だけを持つ行のリストを返します。
判断の根拠: 重複したidを、最後の行で黙って上書きすると、インポートの結果を予測できません。
レビューする誤った変更の断片:
この断片が入った関数の公開契約と比べてみてください。成功ケース1つでは区別できないなら、拒否されるべき入力や、失敗のあとの状態を観測の対象に選びます。
2. 原本のバイト列のフィンガープリントを固定する
source_digest(raw)は、bytesにSHA-256を適用したhex文字列です。JSONを正規化しません。
判断の根拠: 再開は、同じ原本に対してだけ許可する契約です。
レビューする誤った変更の断片:
hashlib.sha256(raw.replace(b" ",b""))
この断片が入った関数の公開契約と比べてみてください。成功ケース1つでは区別できないなら、拒否されるべき入力や、失敗のあとの状態を観測の対象に選びます。
3. 入力とチェックポイントを保存する
init_db(path)は、imports(id TEXT PRIMARY KEY,digest TEXT NOT NULL,next_index INTEGER NOT NULL)とitems(batch TEXT NOT NULL,id TEXT NOT NULL,value INTEGER NOT NULL,PRIMARY KEY(batch,id))を、冪等に作成します。
判断の根拠: 別々のインポート作業の同じ行idは、分離して保存します。
レビューする誤った変更の断片:
CREATE TABLE imports
この断片が入った関数の公開契約と比べてみてください。成功ケース1つでは区別できないなら、拒否されるべき入力や、失敗のあとの状態を観測の対象に選びます。
4. 別の原本で再開できないようにする
begin(path,batch,digest)は、新しい作業ならnext_index=0を保存して0、既存の同じフィンガープリントならnext_index、別のフィンガープリントならValueErrorです。
判断の根拠: 作業idと行番号が同じでも、入力ファイルは違うことがあります。
レビューする誤った変更の断片:
if False:
この断片が入った関数の公開契約と比べてみてください。成功ケース1つでは区別できないなら、拒否されるべき入力や、失敗のあとの状態を観測の対象に選びます。
5. バッチと位置をアトミックにコミットする
apply_chunk(path,batch,start,rows,fault=lambda index:None)は、現在のnext_index==startのときだけ実行し、そうでなければValueErrorです。rowsを順番にitemsに入れ、各挿入のあとにfault(全体のインデックス)を呼びます。すべて成功したら、next_index=start+len(rows)を保存して返します。
判断の根拠: 行1つごとに別々にコミットすると、チェックポイントと行の状態がずれます。
レビューする誤った変更の断片:
db.commit()
fault(start+offset)
この断片が入った関数の公開契約と比べてみてください。成功ケース1つでは区別できないなら、拒否されるべき入力や、失敗のあとの状態を観測の対象に選びます。
6. 現在の位置を取得する
checkpoint(path,batch)は、next_index、ない作業はNoneです。
判断の根拠: 最後に処理を試みた位置ではなく、最後にコミットされた位置を読みます。
レビューする誤った変更の断片:
return 0 if row else None
この断片が入った関数の公開契約と比べてみてください。成功ケース1つでは区別できないなら、拒否されるべき入力や、失敗のあとの状態を観測の対象に選びます。
7. 作業ごとの結果を分離する
values(path,batch)は、そのbatchの(id,value)タプルを、id昇順で返します。
判断の根拠: 別の作業の同じidの行が結果に混ざらないように、batchを条件に置きます。
レビューする誤った変更の断片:
WHERE batch!=? ORDER BY id
この断片が入った関数の公開契約と比べてみてください。成功ケース1つでは区別できないなら、拒否されるべき入力や、失敗のあとの状態を観測の対象に選びます。
8. 途中の失敗のあと、安全に続ける
import_all(path,batch,raw,size=2,fault=lambda index:None)は、boolを除く正のintのsizeを検証し、parse_rows・source_digest・beginを使います。残りの行をsizeずつapply_chunkで処理し、最終的なcheckpointを返します。
判断の根拠: 最初のバッチの成功を保持したまま、2つ目のバッチの失敗後に再開するシナリオを確認します。
レビューする誤った変更の断片:
begin(path,batch,source_digest(raw))
index=0
この断片が入った関数の公開契約と比べてみてください。成功ケース1つでは区別できないなら、拒否されるべき入力や、失敗のあとの状態を観測の対象に選びます。
現場での姿
入力全体をメモリに読み込む、小さなデータ向けのラボです。大容量ファイルをストリーミングでパースするエンジンだと誇張しません。原本の空白だけが違っても、バイトのフィンガープリントが変わるので、再開を拒否するという、保守的な契約です。外部APIの副作用は、このDBトランザクションには入りません。
次のラボですること
8つのステップが、1つの実行可能な成果物につながります。入力行の契約を確認する → 原本のバイト列のフィンガープリントを固定する → 入力とチェックポイントを保存する → 別の原本で再開できないようにする → バッチと位置をアトミックにコミットする → 現在の位置を取得する → 作業ごとの結果を分離する → 途中の失敗のあと、安全に続ける。
各ステップは、関数やファイルが存在するという事実ではなく、実際の戻り値・例外・状態の変化を検査します。正解を見たあとは、わざと境界の比較や後始末のコードを変えて、どの試験が失敗するかを確認してください。前の試験が次のステップでも維持される理由を説明し、このラボが保証しない運用上の条件を1つ書いてみてください。