TT Lab
Get started
Learn Learning paths Courses

Queues and Asynchronous APIs

Where a List Queue Loses Messages

Continue in TT Lab

Summary

BRPOP removes the message from the queue at the same moment it takes it out. If the consumer dies right after taking it out, that message disappears from the world.

Why this was needed

Building a queue with a Redis list takes two lines. You put items in with LPUSH q:jobs "..." and take them out with BRPOP q:jobs 0. It is simple and fast, and for most side projects this is enough.

The problem is in the failure path. The moment BRPOP returns, that message is already gone from Redis. If the consumer dies while processing it, the message is not there even after a restart. Merely restarting a worker during a deployment makes in-progress jobs evaporate.

How it works

The first solution is LMOVE (RPOPLPUSH in older versions). It atomically moves the item to a "processing" list at the same time as taking it out. After finishing processing, the consumer removes it from the processing list. If it dies, the message stays in the processing list, and a separate reclaimer returns old entries to the original queue. This is a hand-made version of the visibility timeout.

The second solution is Redis Stream. A stream was built with this problem in mind from the start. You add with XADD, create a consumer group, and read with XREADGROUP. A message that was read is not deleted but goes into the PEL (Pending Entries List). When processing finishes, you remove it with XACK. If the consumer dies, it stays in the PEL, and another consumer can check it with XPENDING and take it with XCLAIM.

The three properties are summarized as follows.

Property List Stream
Retention after consumption None (deleted immediately) Kept in the PEL, deleted by XACK
Multiple consumer groups Not possible (once taken out, it is over) Independent consumption per group
Reclaiming failed processing Implement it yourself XPENDING / XCLAIM
Memory Small Large (keeps history, needs MAXLEN)

What you meet in the field

How far is ordering guaranteed? With a single list and a single consumer, it is FIFO. The moment you add consumers, the ordering guarantee is gone. This is because when two consumers take messages A and B respectively, you cannot know which finishes first. What needs ordering is usually at the level of a specific entity, so the practical solution is to shard the queue by entity ID and have one consumer per shard.

And when using a stream, you must not forget MAXLEN. Even after XACK, the entry itself remains in a stream. Unless you set a cap like XADD q:orders MAXLEN ~ 100000 * ..., it keeps eating memory. The tilde (~) means approximate trimming, which is much cheaper.

Three pieces for getting a delivery guarantee

"Never losing a message" is not a single setting; all three points must be in place.

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

① Producer acknowledgement. The publish call returning does not mean the broker received it. You must wait for the broker's acknowledgement (ack). If you send without acknowledgement (fire-and-forget), it disappears the moment the broker dies.

② Broker persistence. If it is only in memory, it disappears on restart. Write it to disk and, if possible, confirm it down to the replicas. Kafka's acks=all and RabbitMQ's durable + persistent are this.

③ The consumer's acknowledgement after processing. If you delete as soon as you take it out, it is lost if the consumer dies during processing. You must acknowledge after processing is finished. This is why BRPOP on a Redis List is dangerous.

If even one of the three points is missing, you lose messages no matter how sturdy the others are.

There is no exactly-once

In a distributed system, "exactly-once delivery" cannot be obtained. What you can obtain is at-least-once delivery + idempotent processing, and that combination looks like exactly-once.

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

So making the consumer idempotent is the only answer. You record the IDs of messages already processed and skip them if the same one arrives.

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

This record also grows without limit, so you set a retention period. Make it longer than the maximum period during which a retry can occur (usually a few days).

Preserve order only where order is needed

To preserve global order, you must give up parallel processing. In most cases per-key ordering is enough.

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

If you choose the key badly, everything piles onto one partition (a hot partition). Check the distribution of values first, and if one value takes a large share of the total, do not split by that key.

What you will do in the next lab

You build a queue with a list to confirm FIFO, reproduce the point where messages are lost, block it with LMOVE, and then move to a stream to work with consumer groups and the PEL yourself. At the end you write a table comparing the two on your own.