TT Lab
Get started
Learn Learning paths Courses

Queues and Asynchronous APIs

Implementing a Queue With Redis Lists and Streams

Continue in TT Lab

Goal

Build a queue with a Redis list and with a stream, reproduce the point where the list queue loses messages, and block it with the stream's consumer group and PEL.

Why it matters

Two lines, LPUSH / BRPOP, make a queue. That is why many teams stop here and discover months later that in-progress jobs silently vanish at every deployment. This is because the moment BRPOP returns, the message has already been deleted from Redis. This lab first lets you see that loss with your own eyes and then attaches two solutions in order. One is the hand-made approach of keeping a processing list with LMOVE, and the other is the stream's consumer group, designed for this problem from the start. Once you know the difference between the two, even when you first see SQS's visibility timeout or Kafka's offset commit, you will immediately see what problem the mechanism solves.

Steps

  1. Save the result of redis-cli PING to /root/q/ping.txt. It must contain PONG.
  2. With /root/q/produce.py, put into q:jobs 5 items, from job-1 to job-5. LLEN q:jobs is 5.
  3. With /root/q/consume.py, take out all 5 items and write them one per line to /root/q/order.out. The first line must be job-1 and the last line job-5.
  4. /root/q/safe_consume.py atomically moves items with LMOVE from q:jobs to q:jobs:processing, processes them, and on success removes them from the processing list. Simulate a crash in the middle of processing so that 1 item remains in q:jobs:processing.
  5. With /root/q/stream_add.py, add 5 entries to the stream q:orders using XADD. Apply a cap of MAXLEN ~ 1000. XLEN q:orders is 5.
  6. Create the consumer group g1, read 5 entries with /root/q/stream_consume.py, and XACK only 4 of them.
  7. Save the result of XPENDING q:orders g1 to /root/q/pending.txt. There must be exactly 1 unacknowledged message.
  8. In /root/q/compare.md, write a markdown table. The row titles in the first column are four values, 소비 후 보존, 다중 소비자 그룹, 실패 회수, and 메모리 (retention after consumption, multiple consumer groups, failure reclaim, and memory), and it must have a list column and a stream column.

Notes

Check the Redis connection

Save the result of redis-cli PING to /root/q/ping.txt. It must contain PONG.

Check the response with redis-cli and save the result to a file. It is already running on 127.0.0.1:6379.

Put items into the queue with a list

With /root/q/produce.py, put into q:jobs 5 items, from job-1 to job-5. LLEN q:jobs is 5.

You must put in at one end and take out at the opposite end to get FIFO. Which end you put in at matters.

Prove that the consumption order is FIFO

With /root/q/consume.py, take out all 5 items and write them one per line to /root/q/order.out. The first line must be job-1 and the last line job-5.

Save the order you put in and the order you took out to separate files and compare them. If the order is reversed, the put-in direction and the take-out direction are the same end.

Introduce a processing list to prevent loss

/root/q/safe_consume.py atomically moves items with LMOVE from q:jobs to q:jobs:processing, processes them, and on success removes them from the processing list. Simulate a crash in the middle of processing so that 1 item remains in q:jobs:processing.

Taking out and moving must be a single command to be atomic. If you split it into two commands, the consumer can die in between.

Add messages to the stream

With /root/q/stream_add.py, add 5 entries to the stream q:orders using XADD. Apply a cap of MAXLEN ~ 1000. XLEN q:orders is 5.

Entries are stored as field-value pairs. Also specify a cap so that it does not grow without limit.

Consume with a consumer group and send acknowledgements

Create the consumer group g1, read 5 entries with /root/q/stream_consume.py, and XACK only 4 of them.

You must create the group first to be able to read. If you only read, entries stay in the PEL, and they are removed only when you send an acknowledgement.

Check the messages without an acknowledgement

Save the result of XPENDING q:orders g1 to /root/q/pending.txt. There must be exactly 1 unacknowledged message.

Just deliberately leave one without an acknowledgement. There is a command that queries the number of unacknowledged messages.

Write the list versus stream comparison table

In /root/q/compare.md, write a markdown table. The row titles in the first column are four values, 소비 후 보존, 다중 소비자 그룹, 실패 회수, and 메모리 (retention after consumption, multiple consumer groups, failure reclaim, and memory), and it must have a list column and a stream column.

Organize it along four axes: retention after consumption, multiple groups, failure reclaim, and memory. The table format and the row titles are the grading criteria.