請求はキューで — 失わず、二度処理せず
目標
RabbitMQで保険請求を非同期に渡します。トポロジー(エクスチェンジ・キュー・DLX)を宣言し、発行確認を受け取るパブリッシャーと、手動ack・DLQ・再配信・prefetch・重複排除を備えたコンシューマーを作ります。
なぜ重要なのか
キューは「失わない」ことを約束しますが、その約束は、発行確認と手動ackを行ったときにだけ成り立ち、その代償として同じメッセージが2回届きます。拒否すべきメッセージをキューに戻すとポイズンメッセージになり、prefetchがないと、遅い対象の前でコンシューマーがすべてを抱え込みます。非同期連携の事故の大半は、この5つで起きます。
ステップ
bash /opt/lab/fixtures/eaimw/mq/mq-up.shでブローカーを起動し、rabbitmqctl -n rabbit@localhost status > /root/eaimw/mq/status.txtで状態を残してください。メモリアラームの基準値が絶対値の512MiB(0.5369 gb)になっているか確認します。/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回実行してもエラーが出てはいけません。/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以外のコードで終了します。/root/eaimw/mq/consumer.py --claim <청구시스템URL> --max-seconds <초>(プレースホルダーは請求システムのURLと秒数です)を作成してください。claim.inを消費して本文をPOST /v1/claimsに渡し、201ならそのときにackします(手動ack)。それ以外の結果は、ひとまずキューに戻します(nack、requeue)。--max-secondsが経過したら終了します。- 業務上の拒否(422、書類不備)はキューに戻さず、
reject(requeue=False)でDLXに送ってください。 - 一時エラー(503・接続失敗・タイムアウト)は1回だけキューに戻し、再配信された(
redelivered)メッセージがまた失敗したらDLXに送ってください。 --prefetch(既定は5)でbasic_qosを設定し、遅い請求システムの前でも、ackしていないメッセージが5件を超えないようにしてください。- 処理記録をSQLite(
EAI_MQ_DB、既定は/root/eaimw/mq/processed.db)に残し、すでに処理したmessage_idは請求システムを呼ばずにackだけを行ってください(コンシューマーを再起動しても覚えています)。
参考
- すべてのスクリプトは
EAI_AMQP_URL(既定はamqp://guest:guest@127.0.0.1:5672/%2F)を読みます:pika.BlockingConnection(pika.URLParameters(URL))。採点ツールは一時vhostのアドレスを渡します。guestアカウントは、既定ではループバックからしか接続できません。 - 請求システムのフィクスチャ:
nohup python3 /opt/lab/fixtures/eaimw/partner.py claim > /root/eaimw/mq/claim.out 2>&1 &(9202)。docsが空の請求は422、claimIdがBUSYで始まれば503、統計は/_statsのby_guid(claimIdごとの呼び出し数)です。 - キューの確認:
rabbitmqctl -n rabbit@localhost list_queues name messages messages_unacknowledged。 - pika:
ch.confirm_delivery()のあとでbasic_publish(..., mandatory=True)を呼ぶと、ルーティング不可ならpika.exceptions.UnroutableErrorが出ます。消費はfor m, props, body in ch.consume("claim.in", inactivity_timeout=0.5):です(アイドル中はmがNone)。 - よくある間違い:
auto_ack=Trueにすること、業務上の拒否をrequeueすること、persistentなしでdurableキューだけを作ること。
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の間で死んでも、次に来た同じメッセージは、記録のおかげではじかれます。記録は、コンシューマーを再起動しても残っている必要があります。