Redisのリストとストリームでキューを実装する
目標
Redisのリストとストリームでそれぞれキューを作り、リストのキューがメッセージを失う地点を再現して、ストリームのコンシューマーグループとPELでそれを防ぎます。
なぜ重要なのか
LPUSHとBRPOPの2行でキューになります。そのため、多くのチームがここで止まり、デプロイのたびに処理中だったジョブが静かに消えることを、数か月後に発見します。BRPOPが戻った瞬間に、メッセージはすでにRedisから消えているからです。このラボは、その消失をまず目で確認してから、2つの解決策を順番に付けます。LMOVEで処理中のリストを置く手作りの方法と、最初からこの問題のために設計されたストリームのコンシューマーグループです。この2つの違いがわかれば、SQSの可視性タイムアウトやKafkaのオフセットコミットを初めて見るときも、どんな問題を解く仕組みなのかがすぐに見えます。
ステップ
redis-cli PINGの結果を/root/q/ping.txtに保存します。PONGが入っている必要があります。/root/q/produce.pyでq:jobsにjob-1からjob-5まで5件を入れます。LLEN q:jobsが5です。/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件が残るようにします。/root/q/stream_add.pyで、ストリームq:ordersに5件をXADDします。MAXLEN ~ 1000の上限をかけます。XLEN q:ordersが5です。- コンシューマーグループ
g1を作り、/root/q/stream_consume.pyで5件を読んで、そのうち4件だけをXACKします。 XPENDING q:orders g1の結果を/root/q/pending.txtに保存します。未確認のメッセージがちょうど1件である必要があります。/root/q/compare.mdにMarkdownの表を書きます。1列目の行の見出しは소비 후 보존、다중 소비자 그룹、실패 회수、메모리の4つで(韓国語の4語は、順に「消費後の保持」「複数のコンシューマーグループ」「失敗の回収」「メモリ」を意味します)、リストとストリームの列が必要です。
参考
LMOVE q:jobs q:jobs:processing RIGHT LEFTは、取り出しと移動を1回で行います。- グループの作成:
XGROUP CREATE q:orders g1 0 - 未確認の取得:
XPENDING q:orders g1 - よくある間違い1は、
LPUSHで入れてLPOPで取り出すことです。そうするとLIFOになります。 - よくある間違い2は、ストリームに
MAXLENをかけず、メモリが際限なく伸びることです。
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つの軸で整理します。表の形式と行の見出しが採点基準です。