ファイルキューのコンシューマと冪等な再処理を実装する
目標
ディレクトリベースのキューでコンシューマーを作り、パース失敗を隔離し、DB制約で冪等性を保証し、滞留の監視と再処理まで実装できるようになります。
なぜ重要なのか
非同期連携では、重複は例外ではなくデフォルトです。「exactly-once」は、配信層の魔法ではなく、「少なくとも1回の配信+受信側の重複排除」の結果です。そのため、コンシューマーは常に重複を前提に作る必要があります。そして、パース失敗のメッセージを再試行キューに入れると、そのメッセージが先頭で失敗し続け、後ろの正常なメッセージを塞ぎます(ポイズンメッセージ)。詰まったキューは、そのまま業務の停止です。この2つを自分の手で実装してみると、非同期設計の勘が身につきます。
ステップ
/root/q/inbox、/root/q/processing、/root/q/done、/root/q/errorを作成し、/opt/lab/fixtures/eai/queue/inbox/のすべての.jsonを/root/q/inboxにコピーしてください。inboxに20個のファイルがある必要があります。- メッセージの構造を把握して、
/root/q/schema.csvを作成してください。1行目はfield,type,roleです。正常なメッセージにあるフィールドをすべて書き、role列には、冪等キーの役割をするフィールドにidempotency-keyを、順序の判定に使うフィールドにsequenceを書いてください。 /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は空である必要があります。
- 実行結果が次のようになっている必要があります。
doneが18個、errorが2個processedテーブルの行数は、重複排除後の15件- 重複として無視された
msg_idの3件を、/root/q/dup.txtに昇順で保存
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である必要があります。/root/q/lag.shを作成してください。2つの引数(큐디렉터리 임계치。プレースホルダーはキューのディレクトリとしきい値です)を受け取り、inbox=<n> processing=<n> error=<n>を1行で出力し、inboxの件数がしきい値を超えたら、0以外の終了コードで終了します。/root/q/replay.shを作成してください。引数を1つ(ファイル名)受け取り、errorのそのファイルをinboxに戻します。該当のファイルがerrorになければ、キューの状態をまったく変えずに、0以外の終了コードで終了します。
参考
- sqliteスキーマの確認:
sqlite3 /root/q/ledger.db '.schema processed' - 重複の挿入を無視:
INSERT OR IGNOREまたはON CONFLICT DO NOTHING - アトミックな移動: 同じファイルシステムの中で
os.rename/mv - よくあるミス1: パース失敗のメッセージをinboxに戻して、無限ループを作ってしまうことです。
- よくあるミス2: 重複防止をアプリケーションの条件文だけで行うことです。同時実行されると突破されます。DB制約が最後の防衛線です。
- よくあるミス3:
processingを使わずに、inboxから直接読むことです。コンシューマーを2つ起動すると、両方が同じメッセージを処理します。
キューディレクトリの構成
/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を作成して実行してください。動作は次のとおりです。
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は空である必要があります。
パース失敗は、再試行しても同じように失敗します。そのようなメッセージはすぐにerrorに送らないと、後ろの正常なメッセージが詰まります。処理履歴は、業務処理と同じトランザクションで残してください。
冪等性の確認
実行結果が次のようになっている必要があります。
doneが18個、errorが2個processedテーブルの行数は、重複排除後の15件- 重複として無視された
msg_idの3件を、/root/q/dup.txtに昇順で保存
重複をコードで取り除くより、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を受け取ったら、何もせずに失敗する必要があります。