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

取り消せない変更

宇宙人のお祭りのキャンセルが25%で止まった

TT Labで続きを見る

目標

エイリアン・デザート・フェスティバルの注文40件を、10件ずつキャンセルします。承認一覧は固定し、途中でプロセスが消えても、確定したチャンクの次から再開します。

なぜ重要なのか

全体のロールバックをチャンクごとの確定に変えると、業務上の契約も変わります。部分的な完了を明示的に承認してもらい、チェックポイントが実際の業務上の変更・監査と合っているかを検査する必要があります。先に、Pythonの関数・例外、SQLのトランザクション、前の承認バージョン・補償のラボを学習してください。想定所要時間は120分なので、期限が切れる前に+時間で延長してください。セッションが終わるとファイルが消えます。必要なコードは別に保管してください。

環境と共通契約

成果物は/root/chunks/worker.pyです。PostgreSQL 16・psycopg 3.2.3・Python 3がイメージに用意されていて、実行時のインストールはありません。postgresユーザーで/rootに書き込めるので、追加のcapabilityやユーザーの切り替えは必要ありません。

採点ツールは、ローカルのlabdbの別の一時スキーマに、次のテーブルと仮想の注文を準備し、自分が作ったスキーマだけを片付けます。受講生の関数は、渡された接続のsearch_pathとDSNを使います。publicのテーブルを変更したり、スキーマ・顧客・ID・DSNをハードコードしたりしないでください。SQLの値はパラメータで渡します。

CREATE TABLE orders(id integer PRIMARY KEY,tenant text NOT NULL,
 qty integer NOT NULL CHECK(qty BETWEEN 1 AND 1000),
 state text NOT NULL CHECK(state IN ('pending','paid','cancelled')),
 revision integer NOT NULL CHECK(revision>=0));
CREATE TABLE jobs(job_id text PRIMARY KEY,tenant text NOT NULL,targets jsonb NOT NULL,
 chunk_size integer NOT NULL CHECK(chunk_size BETWEEN 1 AND 10),
 next_index integer NOT NULL CHECK(next_index>=0));
CREATE TABLE job_audit(job_id text REFERENCES jobs(job_id),ordinal integer NOT NULL,
 id integer NOT NULL,previous_revision integer NOT NULL,new_revision integer NOT NULL,
 qty integer NOT NULL,PRIMARY KEY(job_id,ordinal),UNIQUE(job_id,id));

識別子job_id・tenantは厳密なstrで、ASCIIの英数字・アンダースコア・ハイフンの1–64文字です。items/targetsは1–64個の厳密なlistで、各項目はid・revision・qtyだけを持つ厳密なdictです。idはint 1–2147483647、revisionはint 0–2147483646、qtyはint 1–1000です。boolは整数として許可しません。重複したIDは拒否し、ID順のディープコピーに正規化します。chunk_sizeは1–10の厳密なint、allow_partialは正確にTrueでなければなりません。誤った直接の入力は、書き込みの前にValueErrorにし、自動で修正しません。

登録後のjobsのtenant・targets・chunk_sizeと、job_auditは変更しない、という契約です。next_indexは0以上、対象数以下で、完了位置でなければchunk_sizeの倍数でなければなりません。保存されたtargetsは、正規化した入力と同じでなければなりません。監査は、ordinal順に、元の承認配列の[0:next_index]と、ID・前のrevision・新しいrevision=前+1・qtyが正確に一致する必要があります。余った監査も、破損です。元の注文のその後の変更は、過去の監査の破損ではないので、前のチャンクの現在の値で、過去の承認を上書きしません。

run_chunkが処理する区間は[next_index:min(next_index+chunk_size,対象数)]です。faultは、渡されたときだけ、各文字列の引数で呼びます。フックのエラーを隠しません。ロックと文の上限は各SQLごとのもので、チャンク全体の経過時間の制限ではありません。成功・失敗のあとに、元の接続の設定を復元してください。

借りたconはautocommit=True・Read Committedで、外部からの呼び出しの開始時に開いているトランザクションはありません。関数は接続を閉じず、成功・失敗のあとにトランザクションを残しません。cancel_chunkの内部の入れ子の呼び出しは、外側のトランザクションを保持しなければなりません。run_chunkは、実際のコミットのあとにフックを呼ぶので、別の外側のトランザクションで包まないでください。chunk_fileだけが接続を所有します。

ステップ

  1. 部分的な完了の承認を固定します。Exceptionを継承するConflictと、manifest(job_id,tenant,items,chunk_size,allow_partial)を実装してください。下の入力の契約を検証し、job_id・tenant・targets・chunk_size・allow_partialの新しいdictを返します。targetsはID順の新しいlistと新しいdictです。明示的なTrueの承認だけを許可し、誤った入力はValueErrorです。
  2. 再登録しても承認を上書きしません。register(con,plan)は、manifestの5つのキーだけを持つ厳密なdictを検証・正規化したあと、jobsにnext_index=0で登録して、Trueを返します。同じジョブIDで、同じ顧客・正規化したtargets・chunk_sizeなら、何も変更せずFalse、内容が違えばConflictです。入力のエラーは、書き込みの前にValueErrorにし、注文と監査は変更しません。
  3. 承認と進捗を別々に読みます。inspect_job(con,job_id)は、IDを検証し、なければNone、あればjob_id・tenant・targets・chunk_size・next_indexのdictを返します。読み取り専用で、返した値を変更しても元の記録に影響しません。
  4. チャンクの中の1行の失敗でも、すべてロールバックします。cancel_chunk(con,tenant,items)は、対象の入力を先に検証し、ID順に、現在の顧客・ID・revision・qty・pendingを条件にUPDATEします。cancelledに変えてrevisionを1増やし、id・previous_revision・new_revision・qtyのdictのリストを返します。1行でも一致しなければ、Conflictでチャンク全体をロールバックします。jobsとjob_auditには書かず、入れ子の呼び出しは、外側のトランザクションを早期に確定しません。
  5. 変更・監査・位置をまとめて確定します。run_chunk(con,job_id,fault=None)は、自分のトランザクションのロックの上限500ms・文の上限2000msを先に設定し、jobsの行をロックして読みます。行がない場合や、保存された承認・next_index・監査が下の契約と違う場合は、Conflictです。次の区間をcancel_chunkで変更してからfault(after-orders)を呼び、区間の順番・ID・前後のバージョン・数量をすべてjob_auditに保存してからfault(after-audit)を呼びます。next_indexを区間の終わりに変えてからfault(after-checkpoint)を呼び、実際のCOMMITのあとfault(after-commit)を呼びます。戻り値は、processedのID順のlist・next_index・done boolです。すでに完了していれば、processed=[]・done=Trueで、フックは呼びません。コミット前のエラーは、今回のチャンクだけをロールバックし、前のチャンクは保全します。コミット後のエラーは、確定した状態を保持して元のエラーを伝えます。
  6. 過去の完了と現在の違いを報告します。report(con,job_id)は、なければConflict、あればapproved・committed・remaining・matching・drifted・missingのID順のlistを持つdictを返します。approvedは元の承認、committedはnext_indexの前の部分、remainingは残りです。committedの現在の行が、顧客・qty・cancelled・元のrevision+1とすべて同じならmatching、IDがなければmissing、それ以外はdriftedです。進捗と現在の行は1回のSELECTで読み、データは修正しません。
  7. クライアント終了のあとに、残りの区間を引き継ぎます。chunk_file(dsn,job_id,fault=None)は、psycopg.connect(dsn,autocommit=True,connect_timeout=2)で自分の接続を開き、run_chunkを呼んで同じ結果を返します。成功・失敗のどちらでも接続を閉じ、エラーを隠しません。4つのフック地点での実際のクライアント終了のあとの、再開の結果を検査します。
  8. 2つのワーカーで、制限された回数だけ進めます。drain(dsn,job_id,max_chunks)は、IDと、1–64の厳密なintのmax_chunksを、接続の作成の前に検証します。chunk_fileを最大max_chunks回呼びますが、done=Trueならすぐに終わります。呼び出し回数calls・今回の呼び出しが実際に処理したprocessedを連結したリスト・最後のnext_index・最後のdoneを返します。どのエラーもすぐに伝え、隠れた再試行はしません。2つのプロセスが同じジョブを最後まで処理するときに、承認した全IDが、正確に1回ずつだけ変更される必要があります。

参考

部分的な完了の承認を固定する

Exceptionを継承するConflictと、manifest(job_id,tenant,items,chunk_size,allow_partial)を実装してください。下の入力の契約を検証し、job_id・tenant・targets・chunk_size・allow_partialの新しいdictを返してください。targetsはID順の新しいlistと新しいdictにしてください。明示的なTrueの承認だけを許可し、誤った入力はValueErrorにしてください。

40件すべてがアトミックだという約束と、10件ずつ確定する約束は違います。boolとintも区別してください。

再登録しても承認を上書きしない

register(con,plan)は、manifestの5つのキーだけを持つ厳密なdictを検証・正規化したあと、jobsにnext_index=0で登録して、Trueを返してください。同じジョブIDで、同じ顧客・正規化したtargets・chunk_sizeなら、何も変更せずFalse、内容が違えばConflictにしてください。入力のエラーは、書き込みの前にValueErrorにし、注文と監査は変更しないでください。

既存のnext_indexを0で上書きするupsertは、再登録ではなく進捗の消失です。

承認と進捗を別々に読む

inspect_job(con,job_id)は、IDを検証し、なければNone、あればjob_id・tenant・targets・chunk_size・next_indexのdictを返してください。読み取り専用にし、返した値を変更しても元の記録に影響しないようにしてください。

承認配列と現在の注文を混同しないでください。まだ実行していない登録も、有効な記録です。

チャンクの中の1行の失敗でも、すべてロールバックする

cancel_chunk(con,tenant,items)は、対象の入力を先に検証し、ID順に、現在の顧客・ID・revision・qty・pendingを条件にUPDATEしてください。cancelledに変えてrevisionを1増やし、id・previous_revision・new_revision・qtyのdictのリストを返してください。1行でも一致しなければ、Conflictでチャンク全体をロールバックしてください。jobsとjob_auditには書かず、入れ子の呼び出しは、外側のトランザクションを早期に確定しないでください。

UPDATEで影響を受けた行が0でも、SQLは成功です。返る行がないことが業務上の衝突かどうかを判断してください。

変更・監査・位置をまとめて確定する

run_chunk(con,job_id,fault=None)は、自分のトランザクションのロックの上限500ms・文の上限2000msを先に設定し、jobsの行をロックして読んでください。行がない場合や、保存された承認・next_index・監査が下の契約と違う場合は、Conflictにしてください。次の区間をcancel_chunkで変更してからfault(after-orders)を呼び、区間の順番・ID・前後のバージョン・数量をすべてjob_auditに保存してからfault(after-audit)を呼んでください。next_indexを区間の終わりに変えてからfault(after-checkpoint)を呼び、実際のCOMMITのあとfault(after-commit)を呼んでください。戻り値は、processedのID順のlist・next_index・done boolにしてください。すでに完了していれば、processed=[]・done=Trueで、フックは呼ばないでください。コミット前のエラーは、今回のチャンクだけをロールバックし、前のチャンクは保全してください。コミット後のエラーは、確定した状態を保持して元のエラーを伝えてください。

監査の順番を、承認配列の先頭部分の集合と比較してください。同じジョブの行のロックは、次の区間を選ぶ前から必要です。

過去の完了と現在の違いを報告する

report(con,job_id)は、なければConflict、あればapproved・committed・remaining・matching・drifted・missingのID順のlistを持つdictを返してください。approvedは元の承認、committedはnext_indexの前の部分、remainingは残りです。committedの現在の行が、顧客・qty・cancelled・元のrevision+1とすべて同じならmatching、IDがなければmissing、それ以外はdriftedにしてください。進捗と現在の行は1回のSELECTで読み、データは修正しないでください。

現在の別のIDの同じ値は、欠けた元のIDの代わりにはなりません。次のチャンクを進めても、前のチャンクのその後の変更は上書きしません。

クライアント終了のあとに、残りの区間を引き継ぐ

chunk_file(dsn,job_id,fault=None)は、psycopg.connect(dsn,autocommit=True,connect_timeout=2)で自分の接続を開き、run_chunkを呼んで同じ結果を返してください。成功・失敗のどちらでも接続を閉じ、エラーを隠さないでください。4つのフック地点での実際のクライアント終了のあとの、再開の結果を検査します。

コミット後に応答を失ったなら、同じジョブの次の進行は、すでに終わったチャンクではなく次のチャンクです。

2つのワーカーで、制限された回数だけ進める

drain(dsn,job_id,max_chunks)は、IDと、1–64の厳密なintのmax_chunksを、接続の作成の前に検証してください。chunk_fileを最大max_chunks回呼びますが、done=Trueならすぐに終わってください。呼び出し回数calls・今回の呼び出しが実際に処理したprocessedを連結したリスト・最後のnext_index・最後のdoneを返してください。どのエラーもすぐに伝え、隠れた再試行はしないでください。2つのプロセスが同じジョブを最後まで処理するときに、承認した全IDが、正確に1回ずつだけ変更される必要があります。

別々のワーカーが返したIDを、独立したDBの監査と比べます。完了応答を受け取るための空の呼び出しも、callsに含まれます。