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

キューと非同期API

リストキューがメッセージを失う地点

TT Labで続きを見る

一言でいうと

BRPOPは、メッセージをキューから取り出すと同時に消してしまいます。取り出した直後にコンシューマーが死ぬと、そのメッセージはこの世から消えます。

なぜ必要なのか

Redisのリストでキューを作るのは、2行で済みます。LPUSH q:jobs "..."で入れて、BRPOP q:jobs 0で取り出します。簡単で速く、ほとんどのサイドプロジェクトはこれで十分です。

問題は、失敗の経路にあります。BRPOPが戻った瞬間、そのメッセージはすでにRedisから消えています。コンシューマーがそれを処理している途中で死ぬと、再起動してもそのメッセージはありません。デプロイ中にワーカーを再起動するだけで、処理中だったジョブが蒸発します。

どう動くのか

1つ目の解決策はLMOVE(旧バージョンのRPOPLPUSH)です。取り出すと同時に、「処理中」のリストにアトミックに移します。コンシューマーは処理を終えたあと、処理中のリストから消します。死んだ場合、メッセージは処理中のリストに残っており、別の回収プロセスが古い項目を元のキューに戻します。これが可視性タイムアウトの手作り版です。

2つ目の解決策が、Redis Streamです。ストリームは最初からこの問題を念頭に置いて作られました。XADDで追加し、コンシューマーグループを作り、XREADGROUPで読みます。読んだメッセージは消えず、PEL(Pending Entries List)に入ります。処理が終わったらXACKで消します。死んだ場合はPELに残っており、XPENDINGで確認して、XCLAIMでほかのコンシューマーが引き取れます。

3つの性質をまとめると、こうなります。

性質 リスト ストリーム
消費後の保持 なし(即削除) PELに保持、XACKで削除
複数のコンシューマーグループ 不可(一度取り出したら終わり) グループごとに独立して消費
処理失敗の回収 自分で実装 XPENDING / XCLAIM
メモリ 小さい 大きい(履歴を保持、MAXLENが必要)

現場での姿

順序はどこまで保証されるでしょうか。単一のリスト、単一のコンシューマーならFIFOです。コンシューマーを増やした瞬間に、順序の保証はなくなります。2つのコンシューマーがそれぞれメッセージAとBを取っていくと、どちらが先に終わるかわからないからです。順序が必要なのは、たいてい特定のエンティティ単位なので、エンティティIDでキューをシャーディングして、シャードごとにコンシューマーを1つにするのが、実用的な解決策です。

そして、ストリームを使うときにMAXLENを忘れてはいけません。ストリームは、XACKしてもエントリ自体は残ります。XADD q:orders MAXLEN ~ 100000 * ...のように上限を置かないと、メモリを食い続けます。チルダ記号を付けた近似トリミングは、はるかに安上がりです。

配信保証を得る3つの部品

「メッセージを失わない」は、1つの設定ではなく、3つの地点がすべて揃っている必要があります。

생산자 → [브로커] → 소비자
  ①확인      ②지속성    ③확인 후 삭제

① プロデューサーの確認。 発行の呼び出しが戻ったからといって、ブローカーが受け取ったわけではありません。ブローカーの 確認(ack)を待つ必要があります。確認なしで送ると(fire-and-forget)、ブローカーが落ちた瞬間に 消えます。

② ブローカーの永続性。 メモリにしかなければ、再起動で消えます。ディスクに書き、 可能ならレプリカまで確認します。Kafkaのacks=all、RabbitMQのdurableと persistentの組み合わせが、これです。

③ コンシューマーの処理後の確認。 取り出してすぐ削除すると、処理中に死んだときに消えます。 処理が終わったあとに確認する必要があります。Redis ListのBRPOPが危険な理由が、 これです。

3つの地点のうち1つでも欠ければ、残りがどれほど堅牢でも、失います。

ちょうど1回は存在しない

分散システムで「ちょうど1回の配信」は得られません。得られるのは 少なくとも1回の配信+冪等な処理で、その組み合わせがちょうど1回のように見えます。

브로커가 "정확히 한 번" 을 광고하더라도
  → 그것은 브로커 안에서의 이야기다
  → 소비자가 처리하고 확인하기 전에 죽으면 다시 받는다

そのため、コンシューマーを冪等にすることが唯一の答えです。処理したメッセージIDを 記録しておき、同じものが来たら飛ばします。

insert into processed(msg_id) values (%s) on conflict do nothing
-- 삽입된 행이 0이면 이미 처리한 것 → 건너뛴다

この記録も際限なく伸びるので、保持期間を置きます。リトライが起きうる最大の期間 (たいてい数日)より長く設定します。

順序が必要なところだけ順序を守る

グローバルな順序を守るには、並列処理を諦める必要があります。ほとんどはキー単位の順序で 十分です。

파티션 키 = 사용자 ID
  → 같은 사용자의 이벤트는 같은 파티션 → 순서 보장
  → 다른 사용자끼리는 병렬 처리

キーを間違えて選ぶと、1つのパーティションに集中します(ホットパーティション)。値の分布を先に確認し、 1つの値が全体の大きな割合を占めるなら、そのキーでは分けません。

次のラボですること

リストでキューを作ってFIFOを確認し、消失の地点を再現し、LMOVEで防ぎ、そのあとストリームに移してコンシューマーグループとPELを自分で扱います。最後に、2つを比較する表を自分で書きます。