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

システム間連携 (EAI)

ファイルキューのコンシューマと冪等な再処理を実装する

TT Labで続きを見る

目標

ディレクトリベースのキューでコンシューマーを作り、パース失敗を隔離し、DB制約で冪等性を保証し、滞留の監視と再処理まで実装できるようになります。

なぜ重要なのか

非同期連携では、重複は例外ではなくデフォルトです。「exactly-once」は、配信層の魔法ではなく、「少なくとも1回の配信+受信側の重複排除」の結果です。そのため、コンシューマーは常に重複を前提に作る必要があります。そして、パース失敗のメッセージを再試行キューに入れると、そのメッセージが先頭で失敗し続け、後ろの正常なメッセージを塞ぎます(ポイズンメッセージ)。詰まったキューは、そのまま業務の停止です。この2つを自分の手で実装してみると、非同期設計の勘が身につきます。

ステップ

  1. /root/q/inbox、/root/q/processing、/root/q/done、/root/q/errorを作成し、/opt/lab/fixtures/eai/queue/inbox/のすべての.jsonを/root/q/inboxにコピーしてください。inboxに20個のファイルがある必要があります。
  2. メッセージの構造を把握して、/root/q/schema.csvを作成してください。1行目はfield,type,roleです。正常なメッセージにあるフィールドをすべて書き、role列には、冪等キーの役割をするフィールドにidempotency-keyを、順序の判定に使うフィールドにsequenceを書いてください。
  3. /root/q/consume.pyを作成して実行してください。動作は次のとおりです。
    • inboxのファイルを1つずつprocessingに移してから読みます。
    • JSONのパースに失敗した場合、または必須フィールド(msg_id、order_no、seq、amount)が欠落している場合は、errorに移します。
    • 正常なら、sqliteのDB/root/q/ledger.dbのprocessedテーブルにロードして、doneに移します。
    • processedテーブルは、msg_idをPRIMARY KEYまたはUNIQUEとして持つ必要があり、order_no、seq、amount、processed_atのカラムが必要です。
    • 実行が終わったら、processingは空である必要があります。
  4. 実行結果が次のようになっている必要があります。
    • doneが18個、errorが2個
    • processedテーブルの行数は、重複排除後の15件
    • 重複として無視されたmsg_idの3件を、/root/q/dup.txtに昇順で保存
  5. errorに行ったメッセージの理由を、/root/q/error.csvに整理してください。1行目はfile,reasonで、reasonはparseまたはmissing-fieldです。
  6. /root/q/order-check.sqlを作成してください。processedテーブルから、同じorder_noの中でseqが重複している、または欠けている場合を探すクエリです。その結果(問題がなければ0行)を/root/q/order-result.txtに保存してください。問題がなければ、ファイルの1行目はOKである必要があります。
  7. /root/q/lag.shを作成してください。2つの引数(큐디렉터리 임계치。プレースホルダーはキューのディレクトリとしきい値です)を受け取り、inbox=<n> processing=<n> error=<n>を1行で出力し、inboxの件数がしきい値を超えたら、0以外の終了コードで終了します。
  8. /root/q/replay.shを作成してください。引数を1つ(ファイル名)受け取り、errorのそのファイルをinboxに戻します。該当のファイルがerrorになければ、キューの状態をまったく変えずに、0以外の終了コードで終了します。

参考

キューディレクトリの構成

/root/q/inbox、/root/q/processing、/root/q/done、/root/q/errorを作成し、/opt/lab/fixtures/eai/queue/inbox/のすべての.jsonを/root/q/inboxにコピーしてください。 inboxに20個のファイルがある必要があります。

inbox/processing/done/errorの4つの区画を作ります。同じファイルシステムの中に置いて初めて、移動がアトミックになることを覚えておいてください。

メッセージ構造の把握

メッセージの構造を把握して、/root/q/schema.csvを作成してください。 1行目はfield,type,roleです。正常なメッセージにあるフィールドをすべて書き、role列には、冪等キーの役割をするフィールドにidempotency-keyを、順序の判定に使うフィールドにsequenceを書いてください。

メッセージのサンプルを開いて、フィールドを整理します。どのフィールドが冪等キーになれるかに注目してください。

コンシューマーの実装と実行

/root/q/consume.pyを作成して実行してください。動作は次のとおりです。

パース失敗は、再試行しても同じように失敗します。そのようなメッセージはすぐにerrorに送らないと、後ろの正常なメッセージが詰まります。処理履歴は、業務処理と同じトランザクションで残してください。

冪等性の確認

実行結果が次のようになっている必要があります。

重複をコードで取り除くより、DB制約で防ぐほうが安全です。アプリケーションにバグがあっても、制約は破られません。

失敗理由の分類

errorに行ったメッセージの理由を、/root/q/error.csvに整理してください。 1行目はfile,reasonで、reasonはparseまたはmissing-fieldです。

理由を「形式エラー」と「業務エラー」に分けてみてください。前者はソースの修正が必要で、後者はデータを補ったうえで再投入できます。

順序保証の検証

/root/q/order-check.sqlを作成してください。processedテーブルから、同じorder_noの中でseqが重複している、または欠けている場合を探すクエリです。 その結果(問題がなければ0行)を/root/q/order-result.txtに保存してください。 問題がなければ、ファイルの1行目はOKである必要があります。

同じキーのメッセージが順番どおりに処理されたかを、SQLで確認します。処理時刻ではなく、保存された順番号を比較する必要があります。

滞留のモニタリングスクリプト

/root/q/lag.shを作成してください。2つの引数(큐디렉터리 임계치。プレースホルダーはキューのディレクトリとしきい値です)を受け取り、inbox=<n> processing=<n> error=<n>を1行で出力し、inboxの件数がしきい値を超えたら、0以外の終了コードで終了します。

引数でキューのディレクトリを受け取れば、再利用できます。しきい値を超えたら終了コードで知らせて初めて、cronや監視ツールに連携できます。

再処理スクリプト

/root/q/replay.shを作成してください。引数を1つ(ファイル名)受け取り、errorのそのファイルをinboxに戻します。 該当のファイルがerrorになければ、キューの状態をまったく変えずに、0以外の終了コードで終了します。

再処理は、何でも再投入することではありません。存在しないメッセージIDを受け取ったら、何もせずに失敗する必要があります。