distributedsystem systemdesign database/redis

Redis Queue
LMOVE | Docs

Warning

It is now recommended to use Redis Streams for this instead of the recipe described in “Pattern: Reliable Queue”.

A Reliable queue pattern ensures that messages are processed exactly once, even if the consumer crashes. Redis provides the LMOVE command, which moves messages between lists atomically, ensuring message reliability.

How it Works ?

Problem: Redis Lists can lose messages if a consumer crashes after popping a message. Instead, the Reliable Queue Pattern uses two lists.

  1. Main Queue(pending queue) —> Stores new messages
  2. Processing Queue —> Temporary queue for messages being processed.

Instead of RPOP, we use:

LMOVE pending_queue processing_queue RIGHT LEFT

RIGHT → Take the element from the end (right) of pending_queue
LEFT → Insert it at the front (left) of processing_queue

  • LMOVE command return that value to consumer and move it(copy) to another queue(processing_queue).(RPOP only return value)
  • The message moves atomically from pending_queue to processing_queue.
  • If processing fails, unprocessed messages stay in processing_queue for retries.
flowchart TD
A[Producer] -->|LPUSH| B[pending_queue]
B -->|LMOVE RIGHT LEFT| C[processing_queue]
C -->|Consumer Processes| D{Success?}
D -->|YES, LREM| E[Acknowledgment, Remove from processing_queue]
D -->|NO,LMOVE RIGHT LEFT| B[Requeue to pending_queue]

Producer

  • Adds messages to pending_queue using LPUSH.

Consumer

  • Uses LMOVE to atomically transfer messages to processing_queue.
  • Processes the message.
  • If successful, the message is removed.
  • If failed, the message is moved back for retry.

Retry Mechanism

  • A background worker scans processing_queue for stuck messages (e.g., processing timeout).
  • Moves unprocessed messages back to pending_queue using LMOVE.

Go Implementation

package main
 
import (
	"context"
	"fmt"
	"log"
	"time"
 
	"github.com/redis/go-redis/v9"
)
 
var ctx = context.Background()
var rdb = redis.NewClient(&redis.Options{
	Addr: "localhost:6379",
})
 
// Adds a message to the pending queue
func enqueueMessage(message string) error {
	return rdb.LPush(ctx, "pending_queue", message).Err()
}
 
// Moves message from pending queue to processing queue
func dequeueMessage() (string, error) {
	msg, err := rdb.LMove(ctx, "pending_queue", "processing_queue", "RIGHT", "LEFT").Result()
	if err == redis.Nil {
		return "", nil // No messages available
	}
	return msg, err
}
 
// Process messages reliably
func processMessages() {
	for {
		message, err := dequeueMessage()
		if err != nil {
			log.Println("Error moving message:", err)
			continue
		}
 
		if message == "" {
			time.Sleep(1 * time.Second) // No messages, wait before polling
			continue
		}
 
		fmt.Println("Processing message:", message)
 
		// Simulating failure or success
		if time.Now().Unix()%2 == 0 {
			fmt.Println("Message processed successfully:", message)
			rdb.LRem(ctx, "processing_queue", 1, message) // Remove from processing queue
		} else {
			fmt.Println("Processing failed, retrying:", message)
			// Move back to pending queue
			rdb.LMove(ctx, "processing_queue", "pending_queue", "RIGHT", "LEFT")
		}
	}
}
 
func main() {
	// Simulating producer
	enqueueMessage("task-1")
	enqueueMessage("task-2")
	enqueueMessage("task-3")
 
	// Start processing
	processMessages()
}