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

キューと非同期API

Redisのリストとストリームでキューを実装する

TT Labで続きを見る

目標

Redisのリストとストリームでそれぞれキューを作り、リストのキューがメッセージを失う地点を再現して、ストリームのコンシューマーグループとPELでそれを防ぎます。

なぜ重要なのか

LPUSHとBRPOPの2行でキューになります。そのため、多くのチームがここで止まり、デプロイのたびに処理中だったジョブが静かに消えることを、数か月後に発見します。BRPOPが戻った瞬間に、メッセージはすでにRedisから消えているからです。このラボは、その消失をまず目で確認してから、2つの解決策を順番に付けます。LMOVEで処理中のリストを置く手作りの方法と、最初からこの問題のために設計されたストリームのコンシューマーグループです。この2つの違いがわかれば、SQSの可視性タイムアウトやKafkaのオフセットコミットを初めて見るときも、どんな問題を解く仕組みなのかがすぐに見えます。

ステップ

  1. redis-cli PINGの結果を/root/q/ping.txtに保存します。PONGが入っている必要があります。
  2. /root/q/produce.pyでq:jobsにjob-1からjob-5まで5件を入れます。LLEN q:jobsが5です。
  3. /root/q/consume.pyで5件をすべて取り出して、/root/q/order.outに1行ずつ書きます。最初の行がjob-1、最後の行がjob-5である必要があります。
  4. /root/q/safe_consume.pyは、LMOVEでq:jobsからq:jobs:processingへアトミックに移してから処理し、成功したら処理中のリストから消します。処理の途中で死ぬ場合をまねて、q:jobs:processingに1件が残るようにします。
  5. /root/q/stream_add.pyで、ストリームq:ordersに5件をXADDします。MAXLEN ~ 1000の上限をかけます。XLEN q:ordersが5です。
  6. コンシューマーグループg1を作り、/root/q/stream_consume.pyで5件を読んで、そのうち4件だけをXACKします。
  7. XPENDING q:orders g1の結果を/root/q/pending.txtに保存します。未確認のメッセージがちょうど1件である必要があります。
  8. /root/q/compare.mdにMarkdownの表を書きます。1列目の行の見出しは소비 후 보존、다중 소비자 그룹、실패 회수、메모리の4つで(韓国語の4語は、順に「消費後の保持」「複数のコンシューマーグループ」「失敗の回収」「メモリ」を意味します)、リストとストリームの列が必要です。

参考

Redisへの接続を確認する

redis-cli PINGの結果を/root/q/ping.txtに保存してください。PONGが入っている必要があります。

redis-cliで応答を確認して、結果をファイルに残します。127.0.0.1:6379ですでに起動しています。

リストでキューに入れる

/root/q/produce.pyでq:jobsにjob-1からjob-5まで5件を入れてください。LLEN q:jobsが5です。

一方の端に入れて、反対側の端から取り出すとFIFOになります。どちら側に入れるかが重要です。

消費の順序がFIFOであることを証明する

/root/q/consume.pyで5件をすべて取り出して、/root/q/order.outに1行ずつ書いてください。最初の行がjob-1、最後の行がjob-5である必要があります。

入れた順序と取り出した順序をそれぞれファイルに残して、比較してください。順序が逆になっているなら、入れる方向と取り出す方向が同じ側です。

消失を防ぐ処理中リストを導入する

/root/q/safe_consume.pyは、LMOVEでq:jobsからq:jobs:processingへアトミックに移してから処理し、成功したら処理中のリストから消してください。処理の途中で死ぬ場合をまねて、q:jobs:processingに1件が残るようにします。

取り出しと移動を1つのコマンドで行って、初めてアトミックになります。2つのコマンドに分けると、その間に死ぬことがあります。

ストリームにメッセージを追加する

/root/q/stream_add.pyで、ストリームq:ordersに5件をXADDしてください。MAXLEN ~ 1000の上限をかけます。XLEN q:ordersが5です。

フィールドと値のペアで保存されます。際限なく伸びないように、上限も一緒に指定してください。

コンシューマーグループで消費して確認応答を送る

コンシューマーグループg1を作り、/root/q/stream_consume.pyで5件を読んで、そのうち4件だけをXACKしてください。

グループを先に作ってはじめて読めます。読むだけならPELに残り、確認応答を送ると消えます。

確認応答のないメッセージを確認する

XPENDING q:orders g1の結果を/root/q/pending.txtに保存してください。未確認のメッセージがちょうど1件である必要があります。

わざと1つだけ確認応答を送らなければ済みます。未確認のメッセージ数を取得するコマンドがあります。

リストとストリームの比較表を書く

/root/q/compare.mdにMarkdownの表を書いてください。1列目の行の見出しは소비 후 보존、다중 소비자 그룹、실패 회수、메모리の4つで(韓国語の4語は、順に「消費後の保持」「複数のコンシューマーグループ」「失敗の回収」「メモリ」を意味します)、リストとストリームの列が必要です。

消費後の保持、複数グループ、失敗の回収、メモリの4つの軸で整理します。表の形式と行の見出しが採点基準です。