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

EAI 中間層をつくる

キューは失わないが二度届ける

TT Labで続きを見る

一言でいうと

非同期連携は、「今すぐ処理してください」ではなく「失わずにあとで処理してください」を約束します。キューはその約束を守る仕組みですが、約束が成り立つには、パブリッシャー側はブローカーが受け取ったという確認(publisher confirms)を、コンシューマー側は処理したあとにだけ送る確認(手動ack)を行う必要があり、その代償として同じメッセージが2回届くことがあります。

なぜ必要なのか

保険請求の受付を考えてみましょう。顧客がアプリで請求書をアップロードすると、チャネルはすぐに「受け付けました」と表示したくなります。ところが請求システムは書類の検証のために1件あたり数秒かかり、月末には数時間も滞留します。同期で呼び出すとチャネルが請求システムの速度に縛られ、請求システムが少し止まっただけで受付そのものが失敗します。キューを間に置けば、チャネルはキューに入れた瞬間に答えられ、請求システムは自分のペースで取り出していきます。片方が止まっても、もう片方は働き続けます。

その代わり、新しい疑問が生まれます。入れたことをどうやって知るのか。取り出した側が処理の途中で死んだら、メッセージはどこへ行くのか。処理できないメッセージ(書類不備)は、永遠にキューを回り続けるのか。請求システムが遅いと、コンシューマーは何件まで抱え込むのか。このモジュールは、これらの疑問に1つずつ答えます。

どう動くのか

AMQP 0-9-1の3つの部品。パブリッシャーはキューではなくエクスチェンジ(exchange)に送ります。エクスチェンジは、ルーティングキーとバインディングを見てメッセージをキューに入れます。directエクスチェンジは、ルーティングキーがバインディングキーと完全に一致するキューへ送ります(AMQPの概念)。パブリッシャーがキュー名を知らなくてよいこの一段のおかげで、あとから監査用のキューを1つ追加でバインドしても、パブリッシャーは直すところがありません。

durableとpersistentは違います。durableキューは、ブローカーが再起動してもキューの定義が残ります。メッセージが残るには、発行時にpersistent(delivery_mode=2)で送る必要があります。どちらか一方だけだと、再起動後にキューはあるのに空になっています。

パブリッシャー確認(publisher confirms)。basic_publishがエラーなしで戻ってきたからといって、ブローカーが受け取ったわけではありません。AMQPチャネルを確認モードにすると、ブローカーがメッセージごとにackを返し、ドキュメントによれば、durableキューへ向かうpersistentメッセージはディスクに書き込んだあとに確認します(Confirms)。ルーティング先のないメッセージは黙って捨てられますが、mandatoryを付けて送れば、ブローカーがackより先にbasic.returnで返します。pikaのBlockingChannelは、これをUnroutableErrorとして知らせます。

手動ack。自動ackモードでは、ブローカーがメッセージを送った瞬間に配信が終わったものとみなすので、コンシューマーが処理の途中で死ぬと、そのメッセージは消えます。手動ackでは、コンシューマーが処理を終えてackを送って初めて完了です。ackしないままAMQPチャネルが閉じられると、ブローカーはそのメッセージを自動的にキューへ戻し、再び渡すときにredeliveredの印を付けます(同じドキュメント)。これが少なくとも1回(at-least-once)の配信で、裏返して言えば、処理は終わったのにack直前に死んだメッセージがもう一度届きます。キューは冪等を与えてくれません。コンシューマーがmessage_idで処理記録を残して、はじく必要があります(モジュール8の元帳と同じ考え方です)。

拒否には2種類あります。書類不備(業務上の拒否)は、100回渡し直しても結果が同じです。これをキューに戻すと(requeue)、同じメッセージがキューを無限に回ります。一般にポイズンメッセージ(poison message)と呼びます。basic.reject(またはnack)にrequeue=Falseを指定すると、ブローカーはメッセージを捨てるか、キューにデッドレターエクスチェンジ(DLX)が指定されていればそこへ再発行します。キュー引数x-dead-letter-exchange・x-dead-letter-routing-keyで指定し、再発行されたメッセージには、x-deathヘッダーに理由(rejected・expired・maxlen・delivery_limit)が残ります(DLX)。一方、一時エラー(請求システムの503)は、少しあとにやり直せばよいものです。このラボでは1回だけ戻し、redeliveredなのにまた失敗したらDLQへ送ります。(ドキュメントは、本番では引数ではなくポリシー(policy)でDLXを設定することを勧めています。再デプロイせずに変更できるからです。)

prefetchはバックプレッシャーです。コンシューマーのAMQPチャネルのbasic.qos(prefetch_count)は、ackしないまま抱え込める最大件数です。0(無制限)だと、ブローカーはキューのメッセージをすべてコンシューマーに押し込みます。請求システムが遅いと、コンシューマーのメモリに数千件が溜まり、コンシューマーをもう1つ起動しても、すでに全部持っていかれたあとなので分けるものがありません。上限を設ければ、残りはキューに残り、新しいコンシューマーが分け合います。

ブローカーも自分を守ります。メモリがアラームの基準値を超えると、ブローカーは発行する接続をすべてブロックし、消費が進んでメモリが下がれば解除します(メモリアラーム)。基準値の既定は、検出したメモリに対する比率(現在のドキュメントでは0.6)ですが、ドキュメントは、コンテナではブローカーがcgroupの上限を常に把握できるとは限らないと警告し、絶対値を勧めています。実際にこの作成過程で測定したところ、8GBのマシンの2Giコンテナの中で、Ubuntuパッケージの3.12のブローカーは基準値を3.3GBに設定しました。アラームが鳴る前にコンテナがOOMで死ぬということです。そのためこのラボのヘルパーは、vm_memory_high_watermark.absolute = 512MiBで起動します。

現場での姿

最もよくある事故は、自動ackで作ったコンシューマーがデプロイ中に再起動し、処理中だったメッセージを失うことです。ログには何も残りません。2つ目はポイズンメッセージです。形式が誤ったメッセージ1つが無限に再配信され、コンシューマーのCPUを食い、そのあとの正常なメッセージが数時間滞留します。3つ目は、「MQに入れたから安全だ」という思い込みで、パブリッシャー確認をしないことです。ブローカーがメモリアラームで発行を止めているのに、パブリッシャーはタイムアウトだけを見てリトライし、そのあいだに何が入ったのか誰にもわかりません。4つ目は、DLQを作っておきながら誰も見ないことです。DLQには、監視と再処理の手順がセットで必要です。

次のラボですること

Pod内でRabbitMQを起動し(ヘルパー)、エクスチェンジ・キュー・DLXのトポロジーを宣言し、発行確認を受け取るパブリッシャーを作ります。そのあとコンシューマーを順に育てます。手動ack、業務上の拒否のDLQ、一時エラーの1回だけの再配信、prefetch、message_idによる重複排除の順です。採点ツールは、受講生のキューに触れないよう一時vhostを作り、スクリプトをそこで実行します(そのため、すべてのスクリプトがEAI_AMQP_URLを読みます)。