ACK直後に停止した送信処理を復旧する
一言でいうと
未完了の項目は、確認できるACKを受け取ってからだけ完了印を付け、失敗したときは、すでに終わった項目を保存したまま、同じ業務IDでもう一度送ります。
なぜ必要なのか
受付DBと倉庫DBをそれぞれ安全にしても、送信ループの1行が復旧を台無しにすることがあります。リクエストを送る前にsent=1を保存すると、送信の途中で落ちた項目が、次の照会で消えてしまいます。送ったあとで印を付けると、倉庫への反映と印の間で終了したときに、再送が起こります。どちらも、送信を1回だけにしてくれるわけではありません。このラボは、再送を許可して、受信側の効果を保護する方を選びます。
運用者は、「このリクエストは成功しましたか」よりも、「どの状態まで証拠がありますか」を問う必要があります。ordersにある、outboxに未完了としてある、倉庫のinboxにある、パブリッシャーにACKが戻ってきた、sentが保存された、というのはそれぞれ別の観測です。1つの段階が、次の段階を自動的に証明するわけではありません。障害対応でこの違いを表に整理すれば、不確実性を消すためにデータを削除するミスを減らせます。
どう動くのか
小さなバッチを順番に処理する
pendingは、sent=0の項目を、outboxの挿入順seqで読みます。IDの辞書順や、秒単位の時刻では並べ替えません。IDのzがaより先に受け付けられることもあり、同じ時刻の項目が複数あることもあります。1回に1–16個だけを返し、照会自体は状態を変えません。照会してすぐに完了へ変えると、まだ送信していない項目を失います。
dispatchは、選んだバッチの各項目にsendを呼び出します。正確なTrueのACKを受け取ったら、mark_sentを呼び出します。最初の失敗で止まり、元のエラーを伝えます。先に終わった項目はそのまま完了で、失敗した項目と、まだ試していない後ろの項目は未完了のまま残ります。再実行すると、その地点から続きます。このラボは、1つの送信ループの順序を扱い、複数のワーカーが同じバッチを同時に取得するためのclaim・lease・fencingは実装しません。
例えば、A・B・CのうちAのACKのあとで完了印を保存し、Bの応答が消えたなら、次のバッチにはBとCが残ります。Bがすでに倉庫に適用されていても、同じIDで再送するので、inboxが重複した効果を防ぎます。Cはまだ送っていません。Bの失敗を飛ばして、Cから完了させるポリシーも、システムによっては可能ですが、順序が重要な業務なら別の意味になります。このラボのstop-on-first-failureは、明示的に選んだ契約です。
不要なロックと引数の変更を防ぐ
送信の前にpendingの照会を終えて、SQLiteの書き込みトランザクションを閉じます。sendが遅かったり、応答を受け取れなかったりするあいだ、新しい注文の受付までロックしないためです。テストでは、sendコールバックの中で、別のDB接続がBEGIN IMMEDIATEを開始できるかを確認します。単にcon.in_transactionの値だけを読むよりも、実際に競合する接続が書き込みの境界に入れるかを見るのです。
コールバックに渡したdictを、相手が書き換えられる点も考える必要があります。選んだ項目がBだったのに、コールバックが引数のidをCに書き換えると、その書き換えられた値をそのまま完了印に使って、Cを失うことがあります。教育用のeventには文字列と整数しかないので、dictのコピーで十分です。元の選択値で印を残し、コピーで送信します。入れ子のオブジェクトが入る別の契約なら、浅いコピーだけでは保護されないことがあります。
実際の障害の順序を強制的に作る
最終検査は、一時的なsource DBにsnack-7を入れて、別プロセスでdispatchを実行します。HTTPサーバーは、別のsink DBに受信IDと在庫をコミットします。パブリッシャーは、成功ACKを確認した直後に、after-deliveryフックでos._exit(73)により終了します。finallyや正常なcloseは呼び出されません。最初のプロセスが実際にその地点で終了したかを、終了コードで確認します。
親の採点ツールは、sourceには未完了の項目が残り、sinkには在庫7が残っているかを、独立した接続で読みます。続いて、新しい発行プロセスを開始します。サーバーが観測したリクエストは同じIDの2件で、最終的な在庫は7でなければなりません。2つ目のプロセスで、outboxが完了に変わっていなければなりません。Pythonオブジェクトの状態を作り直すだけでなく、別々のプロセスが同じファイルを開き直すテストです。
別のケースでは、サーバーが受信のコミット後にHTTP応答を送らず、接続を閉じます。この場合、パブリッシャーはエラーを受け取って、項目を未完了のまま残す必要があります。再接続して同じ項目を送ると、サーバーは有効な重複としてACKを返し、在庫はそのままです。誤ったIDのACK、過大な応答、リダイレクトも、完了の証拠として受け入れません。send_httpは、このラボのloopbackの/eventsだけを許可し、プロキシ環境変数を使わず、リダイレクトにも従いません。
現場での姿
実際の運用では、送信成功の回数だけで状態を判断しません。未完了の項目数、最も古い未完了の経過時間、試行回数、受信の重複数、衝突数を分けて見ます。前の項目の1つが永続的な形式エラーなら、後ろの項目がずっと塞がれることがあります。無限リトライの代わりに、アラート・隔離・修正・再投入の手順を用意する必要があります。隔離したからといって成功したわけではなく、順序を入れ替える影響もあわせて記録する必要があります。
リトライの待機には上限とジッターが必要ですが、今回のdispatchは、呼び出し1回の有限のバッチだけを実行します。自動リトライのスケジューラーや、高可用性のあるジョブ分配を実装したとは説明しません。sourceとsinkのファイルも、LabHubのセッションが終わると消えます。ここで検証したのは、プロセスの終了後に同じディスクを再び使う場合です。ホストの電源遮断、ディスク障害、バックアップからの復旧、インターネット上での性能は、別の検証です。
次のラボですること
worker.pyに、イベントの検証、ストアの初期化、原子的な受付、未完了の照会、重複受信の処理、完了印、有限バッチの送信、HTTP ACKの確認を、順番に実装します。提供された採点ツールは、一時DBとloopbackサーバーを自分で作り、正解コードを実行します。空の関数の枠は文法が正しいだけで、答えの代わりにはなりません。ラボのあとのクイズでは、2回の配信・1回の適用という結果が立証する範囲と、残る運用上の責任を整理します。
公式ドキュメントでさらに読む
- Python urllib.request: リクエスト・応答とリダイレクトハンドラーを確認します。
- SQLite Atomic Commit: プロセス終了のテストと、ストレージ装置の耐久性の前提を区別します。