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

ACKの前に止まったお菓子の自販機

二人のお菓子配達員と終わらない再試行

TT Labで続きを見る

目標

学園祭のお菓子の配達員2人が競い合い、送信の直後に1人が落ちても、残った業務を有限のバジェットで復旧します。古いワーカーが、新しいワーカーの完了記録を上書きできないようにします。

なぜ重要なのか

ネットワークのタイムアウトは、業務が処理されていないという意味ではありません。リースが終わっても、以前のリクエストは実行され続けることがあります。今回は、キューの先取りトークンと受信inboxを区別し、ロックの外での送信・失敗バジェット・隔離記録を、あわせて設計します。前のoutboxラボの、重複した効果の防止を、複数の配達員がいる状況へ拡張します。

予想の所要時間は110分です。基本のセッションより長いので、期限が切れる前に+時間を押して延長してください。セッションが終了するとファイルが消えます。必要なコードは別に保管してください。Pythonの例外処理とSQLiteトランザクション、前のモジュールのoutbox・受信の重複排除を理解してから始めます。

データ契約

成果物は/root/lease/worker.pyです。キューのファイルとHTTP受信サーバーは、採点ツールが一時的に作って後始末するので、パス・ポートをハードコードしません。すべてのキューは、正常なスキーマを持つ、信頼したローカルファイルです。下のスキーマを使います。

CREATE TABLE jobs (seq INTEGER PRIMARY KEY AUTOINCREMENT,
  id TEXT NOT NULL UNIQUE, qty INTEGER NOT NULL,
  status TEXT NOT NULL CHECK(status IN ('pending','leased','done','dead')),
  attempts INTEGER NOT NULL, max_attempts INTEGER NOT NULL, token INTEGER NOT NULL,
  available_at INTEGER NOT NULL, deadline INTEGER NOT NULL,
  owner TEXT, lease_until INTEGER, reason TEXT);
CREATE TABLE redrives (seq INTEGER PRIMARY KEY AUTOINCREMENT,
  id TEXT NOT NULL, token INTEGER NOT NULL, at_ms INTEGER NOT NULL, note TEXT NOT NULL);

識別子は、ASCIIの英数字・アンダースコア・ハイフンの1–64文字の、正確なstrです。now_msとclockの戻り値は、共通の時間基準の、正確なintで0から(10**15-60000)までで、すべての整数の契約は、boolを拒否します。個々のDB関数は、形式・範囲のエラーを、変更の前にValueErrorで拒否し、入力を変更しません。run_onceの2つ目のclockが不正だったり、後ろへ戻っていたりした場合は、すでにコミットした先取りを保存して、エラーを伝えます。ホスト間の時計を合わせる合意アルゴリズムではありません。

貸したconは、呼び出しの開始時にトランザクションがなく、関数はこれを閉じません。成功・失敗のあとに、開いたトランザクションを残しません。書き込みの選択と変更は、BEGIN IMMEDIATEで束ね、コミット前のfaultのエラーは全体をロールバックして、元のエラーを伝えます。コミット後のエラーは、すでに確定した状態を保存します。fault=Noneなら呼び出さず、フックは、その変更を行った経路でだけ呼び出します。None・Falseなど変更のない戻りの経路には、フックがありません。claimは、対象がなくても、先に行った隔離の更新をコミットします。

業務ごとに意味が独立しています。前の業務の失敗が、あとの業務を塞ぐ、厳格な順序保証のキューではありません。tokenは、業務ごとの単調増加の先取りの番号で、検査の範囲では10**15以内です。キューの印の検査は、外部サーバーの権限の検証や、外部リソースのfencingの代わりにはなりません。受信サーバーは、同一ID・同一数量の重複した効果を取り除くように提供されます。redriveは、権限のある運用者が呼び出すという前提で、ログイン・ロール管理の機能は実装しません。

ステップ

  1. リトライの遅延の上限とジッターを計算する: worker.pyに、Exceptionのサブクラスとして、Conflict・Retryable・Permanent・BadAckを定義し、retry_delay(attempt,base_ms,cap_ms,jitter)を実装してください。attemptは1–16、base_msは1–60000、cap_msはbase_ms–60000の、正確なintです。jitterはboolを除く正確なint/floatで、有限の0–1です。floor(min(cap_ms,base_ms*2(attempt-1))*jitter)を返します。不正な入力はValueErrorです。内部の乱数やsleepは使わず、与えられた比率で計算してください。
  2. 再起動のあとも残るキューを開く: open_queue(path)は、下のjobs・redrivesテーブルを、存在しないときにだけ1つのトランザクションで作成し、sqlite3.Connectionを返します。isolation_level=None、timeout=1秒、journal_mode=DELETE、synchronous=FULLを使います。既存の業務・監査の行は保存し、初期化に失敗したときは接続を閉じます。インメモリDBではなく、使い捨てのローカルファイルを受け取ります。
  3. 重複した受付でバジェットを初期化しない: enqueue(con,event,now_ms,ttl_ms,max_attempts=3)は、idとqtyだけを持つ正確なdictを受け取ります。idは下の識別子、qtyは正確なintの1–1000、ttl_msは1–60000、max_attemptsは1–16です。新しい業務を、pending、attempts=0、token=0、available_at=now_ms、deadline=now_ms+ttl_ms、owner・lease_until・reason=NULLで挿入し、Trueを返します。同じID・同じ数量は、別の入力バジェットが来ても、既存の行全体を保存してFalse、違う数量はConflictです。先にすべての入力を検証し、選択・挿入を1つの書き込みトランザクションで処理します。
  4. 2人の配達員が同時に同じ仕事を持っていけないようにする: claim(con,owner,now_ms,lease_ms=1000,fault=None)を実装してください。ownerは識別子、lease_msは正確なintの1–60000です。BEGIN IMMEDIATEの中で、pendingまたは期限切れのleasedの行のうち、deadline<=now_msまたはattempts>=max_attemptsの行を、deadに移します。reasonは、締め切りならdeadline、そうでなければexhaustedで、owner・lease_untilはNULLにします。残った行のうち、pendingでavailable_at<=now_ms、またはleasedでlease_until<=now_msの業務を、seqの昇順で1つ選びます。なければNoneです。あれば、attempts・tokenをそれぞれ1上げて、leased、owner、lease_until=min(now_ms+lease_ms,deadline)、reason=NULLで保存します。id・qty・owner・token・attempt・lease_untilだけを持つdictを返します。attemptは、更新したattemptsです。更新のあとにfault('after-claim')、COMMITのあとにfault('after-commit')を呼び出します。
  5. 古い配達員の完了を拒否する: finish(con,ticket,outcome,now_ms,delay_ms=0,fault=None)を実装してください。ticketは、claimの6つのキーだけを持つdictで、id・ownerは識別子、qty=1–1000、token=1–1015、attempt=1–16、lease_until=1–10**15の、正確なintです。outcomeはok・retry・permanentのいずれかのstr、delay_msは正確なintの0–60000です。1つの書き込みトランザクションで、業務がない、leasedではない、ticketの数量・所有者・トークン・試行・期限が違う、now_ms>=lease_untilまたはdeadlineのいずれかなら、変更せずにFalseです。有効なokはdone/理由NULL、permanentはdead/permanentです。retryは、最大試行に達していればdead/exhausted、そうでなければnow_ms+delay_ms>=deadlineならdead/deadline、それ以外はpending/retryです。pendingのときにだけavailable_at=now_ms+delay_ms、それ以外はnow_msで、owner・lease_untilはNULLにします。ほかのフィールドは保存し、新しいstatusの文字列を返します。更新のあとにfault('after-finish')、COMMITのあとにfault('after-commit')を呼び出します。
  6. 隔離の解除と監査記録を一緒に残す: redrive(con,job_id,now_ms,ttl_ms,note,fault=None)は、存在するdeadの業務だけを再投入し、Trueを返します。job_id・noteは識別子で、ttl_msは正確なintの1–60000です。対象がない、またはdeadでなければValueErrorです。1つの書き込みトランザクションで、redrivesにid・現在のtoken・at_ms=now_ms・noteを挿入→fault('after-audit')→業務をpending、attempts=0、available_at=now_ms、deadline=now_ms+ttl_ms、owner・lease_until・reason=NULLに変更→fault('after-redrive')→COMMIT→fault('after-commit')の順です。qty・max_attempts・tokenは保存します。
  7. ネットワーク送信を書き込みロックの外で行う: run_once(con,owner,clock,send,jitter,lease_ms=1000)は、jitterを先に検証し、clock()の最初の時刻でclaimします。業務がなければNoneで、sendは呼び出しません。あれば、トランザクションの外で、idとqtyだけを持つ新しいdictをsendに1回渡します。正確なTrueだけがokで、それ以外の戻り値はBadAckです。Retryableだけをretryに分類して、retry_delay(ticket.attempt,100,1000,jitter)を使い、Permanentはpermanentに分類します。それ以外のエラー・BadAckはそのまま伝え、leased状態を保存します。正常な場合と、分類されたエラーのあとで、clock()をもう一度呼び出し、2つ目の値が最初の値より小さければValueErrorです。finishに2つ目の時刻と結果・遅延を渡してから、id・token・statusのdictを返します。finishがFalseなら、statusはstaleです。送信コールバックが受け取ったdictを変更したり、別の接続で業務を入れたりしても、現在の印は維持されなければなりません。
  8. 送信の直後に死んだ作業を別のプロセスが復旧する: run_file(path,owner,clock,send,jitter,lease_ms=1000)は、既存の通常のキューファイルがなければFileNotFoundErrorで拒否し、空のキューを新しく作りません。open_queueで接続を所有し、run_onceの結果を返し、成功・失敗のどちらでも接続を閉じます。最終検査では、実際の2つの子プロセスの同時claimのうち1つだけが成功するか、コミットの前後の強制終了のあとで状態がアトミックかを検査します。別のHTTP受信サーバーが数量をコミットした直後に、最初の配達員を終了させ、2番目の配達員のプロセスで再送します。HTTPリクエストは2回ですが、inboxと受信の効果は1回でなければなりません。受信サーバーと一時DBは、採点ツールが提供します。

参考

リトライの遅延の上限とジッターを計算する

worker.pyに、Exceptionのサブクラスとして、Conflict・Retryable・Permanent・BadAckを定義し、retry_delay(attempt,base_ms,cap_ms,jitter)を実装してください。attemptは1–16、base_msは1–60000、cap_msはbase_ms–60000の、正確なintです。jitterはboolを除く正確なint/floatで、有限の0–1です。floor(min(cap_ms,base_ms*2**(attempt-1))*jitter)を返します。不正な入力はValueErrorです。内部の乱数やsleepは使わず、与えられた比率で計算してください。

boolはPythonではintのサブタイプです。上限を先に適用し、注入した比率を掛けてから切り捨ててください。

再起動のあとも残るキューを開く

open_queue(path)は、下のjobs・redrivesテーブルを、存在しないときにだけ1つのトランザクションで作成し、sqlite3.Connectionを返します。isolation_level=None、timeout=1秒、journal_mode=DELETE、synchronous=FULLを使います。既存の業務・監査の行は保存し、初期化に失敗したときは接続を閉じます。インメモリDBではなく、使い捨てのローカルファイルを受け取ります。

短い書き込みトランザクションだけでワーカーを調整します。テーブルの作成と、既存の内容の初期化を混同しないでください。

重複した受付でバジェットを初期化しない

enqueue(con,event,now_ms,ttl_ms,max_attempts=3)は、idとqtyだけを持つ正確なdictを受け取ります。idは下の識別子、qtyは正確なintの1–1000、ttl_msは1–60000、max_attemptsは1–16です。新しい業務を、pending、attempts=0、token=0、available_at=now_ms、deadline=now_ms+ttl_ms、owner・lease_until・reason=NULLで挿入し、Trueを返します。同じID・同じ数量は、別の入力バジェットが来ても、既存の行全体を保存してFalse、違う数量はConflictです。先にすべての入力を検証し、選択・挿入を1つの書き込みトランザクションで処理します。

受付のリトライは新しい業務ではありません。UNIQUEのIDと既存の数量を比べ、入力のdictも変更しないでください。

2人の配達員が同時に同じ仕事を持っていけないようにする

claim(con,owner,now_ms,lease_ms=1000,fault=None)を実装してください。ownerは識別子、lease_msは正確なintの1–60000です。BEGIN IMMEDIATEの中で、pendingまたは期限切れのleasedの行のうち、deadline<=now_msまたはattempts>=max_attemptsの行を、deadに移します。reasonは、締め切りならdeadline、そうでなければexhaustedで、owner・lease_untilはNULLにします。残った行のうち、pendingでavailable_at<=now_ms、またはleasedでlease_until<=now_msの業務を、seqの昇順で1つ選びます。なければNoneです。あれば、attempts・tokenをそれぞれ1上げて、leased、owner、lease_until=min(now_ms+lease_ms,deadline)、reason=NULLで保存します。id・qty・owner・token・attempt・lease_untilだけを持つdictを返します。attemptは、更新したattemptsです。更新のあとにfault('after-claim')、COMMITのあとにfault('after-commit')を呼び出します。

選択と更新の間でロックを解かないでください。期限切れの境界の等号と、選ぶ仕事がなくても確定する必要がある隔離の変更を、区別してください。

古い配達員の完了を拒否する

finish(con,ticket,outcome,now_ms,delay_ms=0,fault=None)を実装してください。ticketは、claimの6つのキーだけを持つdictで、id・ownerは識別子、qty=1–1000、token=1–1015、attempt=1–16、lease_until=1–1015の、正確なintです。outcomeはok・retry・permanentのいずれかのstr、delay_msは正確なintの0–60000です。1つの書き込みトランザクションで、業務がない、leasedではない、ticketの数量・所有者・トークン・試行・期限が違う、now_ms>=lease_untilまたはdeadlineのいずれかなら、変更せずにFalseです。有効なokはdone/理由NULL、permanentはdead/permanentです。retryは、最大試行に達していればdead/exhausted、そうでなければnow_ms+delay_ms>=deadlineならdead/deadline、それ以外はpending/retryです。pendingのときにだけavailable_at=now_ms+delay_ms、それ以外はnow_msで、owner・lease_untilはNULLにします。ほかのフィールドは保存し、新しいstatusの文字列を返します。更新のあとにfault('after-finish')、COMMITのあとにfault('after-commit')を呼び出します。

ownerが同じでも、古い実行かもしれません。印全体と時間の境界を確認し、結果が遅れたからといって現在の業務を削除しないでください。

隔離の解除と監査記録を一緒に残す

redrive(con,job_id,now_ms,ttl_ms,note,fault=None)は、存在するdeadの業務だけを再投入し、Trueを返します。job_id・noteは識別子で、ttl_msは正確なintの1–60000です。対象がない、またはdeadでなければValueErrorです。1つの書き込みトランザクションで、redrivesにid・現在のtoken・at_ms=now_ms・noteを挿入→fault('after-audit')→業務をpending、attempts=0、available_at=now_ms、deadline=now_ms+ttl_ms、owner・lease_until・reason=NULLに変更→fault('after-redrive')→COMMIT→fault('after-commit')の順です。qty・max_attempts・tokenは保存します。

試行バジェットは新しく与えますが、先取りの世代は再利用しません。監査の行だけが残ったり、状態だけが解除されたりする、中途半端な再投入を防いでください。

ネットワーク送信を書き込みロックの外で行う

run_once(con,owner,clock,send,jitter,lease_ms=1000)は、jitterを先に検証し、clock()の最初の時刻でclaimします。業務がなければNoneで、sendは呼び出しません。あれば、トランザクションの外で、idとqtyだけを持つ新しいdictをsendに1回渡します。正確なTrueだけがokで、それ以外の戻り値はBadAckです。Retryableだけをretryに分類して、retry_delay(ticket.attempt,100,1000,jitter)を使い、Permanentはpermanentに分類します。それ以外のエラー・BadAckはそのまま伝え、leased状態を保存します。正常な場合と、分類されたエラーのあとで、clock()をもう一度呼び出し、2つ目の値が最初の値より小さければValueErrorです。finishに2つ目の時刻と結果・遅延を渡してから、id・token・statusのdictを返します。finishがFalseなら、statusはstaleです。送信コールバックが受け取ったdictを変更したり、別の接続で業務を入れたりしても、現在の印は維持されなければなりません。

外部の応答を待ちながら、DBトランザクションを握らないでください。2つ目の時刻で有効性をもう一度判断し、わからない失敗を成功として隠さないでください。

送信の直後に死んだ作業を別のプロセスが復旧する

run_file(path,owner,clock,send,jitter,lease_ms=1000)は、既存の通常のキューファイルがなければFileNotFoundErrorで拒否し、空のキューを新しく作りません。open_queueで接続を所有し、run_onceの結果を返し、成功・失敗のどちらでも接続を閉じます。最終検査では、実際の2つの子プロセスの同時claimのうち1つだけが成功するか、コミットの前後の強制終了のあとで状態がアトミックかを検査します。別のHTTP受信サーバーが数量をコミットした直後に、最初の配達員を終了させ、2番目の配達員のプロセスで再送します。HTTPリクエストは2回ですが、inboxと受信の効果は1回でなければなりません。受信サーバーと一時DBは、採点ツールが提供します。

同じファイルを、次のプロセスが開き直します。外部の受信の効果とキューの完了の間のすきまをなくしたと主張せず、重複送信を復旧してください。