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

EAI 中間層をつくる

請求はキューで — 失わず、二度処理せず

TT Labで続きを見る

目標

RabbitMQで保険請求を非同期に渡します。トポロジー(エクスチェンジ・キュー・DLX)を宣言し、発行確認を受け取るパブリッシャーと、手動ack・DLQ・再配信・prefetch・重複排除を備えたコンシューマーを作ります。

なぜ重要なのか

キューは「失わない」ことを約束しますが、その約束は、発行確認と手動ackを行ったときにだけ成り立ち、その代償として同じメッセージが2回届きます。拒否すべきメッセージをキューに戻すとポイズンメッセージになり、prefetchがないと、遅い対象の前でコンシューマーがすべてを抱え込みます。非同期連携の事故の大半は、この5つで起きます。

ステップ

  1. bash /opt/lab/fixtures/eaimw/mq/mq-up.shでブローカーを起動し、rabbitmqctl -n rabbit@localhost status > /root/eaimw/mq/status.txtで状態を残してください。メモリアラームの基準値が絶対値の512MiB(0.5369 gb)になっているか確認します。
  2. /root/eaimw/mq/topology.pyを作成してください。エクスチェンジはeai.claim(direct、durable)とeai.dlx(direct、durable)、キューはclaim.in(durable、引数x-dead-letter-exchange=eai.dlx、x-dead-letter-routing-key=claim.dead)とclaim.dead(durable)、バインディングはeai.claim→claim.in(キーclaim)とeai.dlx→claim.dead(キーclaim.dead)です。2回実行してもエラーが出てはいけません。
  3. /root/eaimw/mq/publish.py <청구JSON파일>(プレースホルダーは請求JSONファイルです)を作成してください。JSONをそのまま本文にして、eai.claimにルーティングキーclaimで発行します。persistent(delivery_mode 2)、message_idはJSONのguid、content_typeはapplication/jsonにします。発行確認を有効にしてmandatoryで送り、ルーティング不可や拒否なら0以外のコードで終了します。
  4. /root/eaimw/mq/consumer.py --claim <청구시스템URL> --max-seconds <초>(プレースホルダーは請求システムのURLと秒数です)を作成してください。claim.inを消費して本文をPOST /v1/claimsに渡し、201ならそのときにackします(手動ack)。それ以外の結果は、ひとまずキューに戻します(nack、requeue)。--max-secondsが経過したら終了します。
  5. 業務上の拒否(422、書類不備)はキューに戻さず、reject(requeue=False)でDLXに送ってください。
  6. 一時エラー(503・接続失敗・タイムアウト)は1回だけキューに戻し、再配信された(redelivered)メッセージがまた失敗したらDLXに送ってください。
  7. --prefetch(既定は5)でbasic_qosを設定し、遅い請求システムの前でも、ackしていないメッセージが5件を超えないようにしてください。
  8. 処理記録をSQLite(EAI_MQ_DB、既定は/root/eaimw/mq/processed.db)に残し、すでに処理したmessage_idは請求システムを呼ばずにackだけを行ってください(コンシューマーを再起動しても覚えています)。

参考

Pod内でブローカーを起動する

mq-up.shでRabbitMQを起動し、rabbitmqctl statusの出力を/root/eaimw/mq/status.txtに残してください。

ヘルパーがノード名をrabbit@localhostで起動するので、rabbitmqctlには-n rabbit@localhostを指定します。出力のMemory high watermarkの行を見てください。

エクスチェンジ・キュー・DLXを宣言する

/root/eaimw/mq/topology.pyが、エクスチェンジ2つ・キュー2つ・バインディング2つと、claim.inのDLX引数を宣言するようにしてください(2回実行しても安全に)。

宣言は冪等です。同じ属性で再宣言しても何も起きず、属性が違うとブローカーが406 PRECONDITION_FAILEDでチャネルを閉じます。DLXはキュー引数2つです。

確認を受け取るパブリッシャー

/root/eaimw/mq/publish.py がpersistent・message_id・content_typeを付けて発行し、発行確認とmandatoryでルーティング不可を検知するようにしてください。

ch.confirm_delivery()を先に呼ぶと、basic_publishがブローカーの確認を待ちます。mandatory=Trueで送ると、受け取るキューがない場合にUnroutableErrorが発生します。

処理したあとにだけackする

/root/eaimw/mq/consumer.pyがclaim.inを消費して請求システムに渡し、201を受け取ったあとにだけackするようにしてください(それ以外はキューに戻す)。

ch.consume(queue, inactivity_timeout=0.5)で回すと、アイドル中はNoneが来ます。--max-secondsが経過したらループを抜けて接続を閉じてください。ackしていないものは、ブローカーがキューに戻します。auto_ackは使いません。

業務上の拒否はDLQへ

422(書類不備)はキューに戻さず、reject(requeue=False)でDLXに送ってください。

戻したメッセージはすぐにまた届きます。100回渡しても結果が同じメッセージを戻すと、コンシューマーがそれだけを抱えて回り続けます。requeue=Falseなら、キューに設定されたDLXへ行きます。

一時エラーはもう1回、それでもだめならDLQ

503・接続失敗は1回だけキューに戻し、redeliveredなのにまた失敗したらDLXに送ってください。

methodフレームのredeliveredがTrueなら、すでに1回戻されたメッセージです。戻すことを無制限に許すと、対象が長く止まっているときに、キュー全体がそのメッセージだけを回り続けます。

prefetchで抱え込む量を制限する

--prefetch(既定は5)でbasic_qosを設定し、unackedが5件を超えないようにしてください。

ch.basic_qos(prefetch_count=N)は、チャネル単位でackしていない配信数を制限します。消費を始める前に呼んでください。

同じmessage_idは1回だけ

処理したmessage_idをSQLite(EAI_MQ_DB)に残し、再び届いたら呼び出さずにackだけを行ってください。

処理に成功したら記録してackします。記録とackの間で死んでも、次に来た同じメッセージは、記録のおかげではじかれます。記録は、コンシューマーを再起動しても残っている必要があります。