952 messages nobody was ever going to read
Redis Streams consumer groups are close enough to Kafka that you can carry the mental model across, and just different enough that the place they differ is the place that hurts.
When a consumer reads from a stream with XREADGROUP, the message is added to
that consumer's pending-entries list — the PEL. It stays there until the
consumer acknowledges it with XACK. This is the delivery guarantee: a message
that was read but never acked is not lost, because Redis still knows it was
delivered and to whom.
That is the feature. Here is the cost: if the consumer that read the message dies, the message stays in that consumer's PEL forever. It is not redistributed. Another consumer in the group will never see it. Redis is faithfully holding it for a process that is not coming back.
What we found
A routine look at stream depth showed the stream itself healthy — trimmed, short, moving. Then:
XPENDING match_events ranking_group
1) (integer) 952
952 messages, the oldest a little over three weeks old, spread across consumer names that no longer corresponded to any running process. Every deploy had rolled the pods, every roll had orphaned whatever those consumers held in-flight, and nothing had ever reclaimed them.
Nothing alerted, because nothing was wrong by any metric being watched. Stream length was fine. Consumer lag — the metric everyone graphs — was fine, because lag measures unread messages and these had all been read. Error rates were zero. The system was silently 952 events behind on a pipeline nobody thought had a backlog.
The missing piece
There was no reclaim path. The code had XREADGROUP and XACK and stopped
there, which is what you write if your model is "Kafka, but Redis". Kafka
rebalances partitions when a consumer disappears; the work moves. Redis
Streams does not. Reclaiming is an explicit operation you have to run
yourself:
// Take ownership of anything idle for longer than the threshold and
// redeliver it to a live consumer.
res, err := rdb.XAutoClaim(ctx, &redis.XAutoClaimArgs{
Stream: stream,
Group: group,
Consumer: myName,
MinIdle: 5 * time.Minute,
Start: "0-0",
Count: 100,
}).Result()
XAUTOCLAIM walks the PEL, finds entries idle longer than MinIdle, and
reassigns them. Run it periodically from every consumer and orphaned work gets
picked up within one interval of the process dying.
Two things to get right that are easy to get wrong:
MinIdlemust be comfortably longer than your slowest legitimate processing time. Set it too low and you will reclaim messages from consumers that are alive and still working, and process them twice for no reason.- A message that keeps failing will be reclaimed forever.
XAUTOCLAIMreports a delivery count; past a threshold the message belongs in a dead letter stream, not back in the rotation. Without that, one poison message becomes an infinite loop that looks like healthy throughput.
The metric that would have caught it
Consumer lag would never have found this. The metric that finds it is oldest-pending-entry age, per group:
XPENDING <stream> <group>
returns the smallest and largest ids in the PEL. The age of the smallest, alarmed above a few minutes, catches an orphaned PEL the day it appears rather than three weeks later.
It is a good example of a general shape: the metric everyone graphs is the one the documentation puts first, and the failure that actually happens is often in the state the documentation mentions in passing.
The full case study is the match service.