Skip to main content

Command Palette

Search for a command to run...

How to Fix Redis Pub/Sub Message Loss

Learn: How to Fix Redis Pub/Sub Message Loss

Updated
5 min readView as Markdown
T

Welcome to TopperBlog! 👋

I'm a tech content creator passionate about helping developers level up their careers and master cutting-edge technologies.

🎯 What I Write About: • AI/ML Engineering & LLMs • Web3 & Blockchain Development
• System Design & Architecture • Interview Preparation (FAANG) • Freelancing & Remote Work • Modern Tech Stacks (Next.js, React, Rust, TypeScript) • Performance Optimization & Best Practices

💼 Mission: Sharing practical, actionable insights that accelerate your tech career and maximize your earning potential.

📚 15+ In-Depth Guides covering everything from earning $10k/month as a freelancer to cracking FAANG interviews.

🌐 Let's connect and grow together in this amazing tech journey!

#TechBlogger #SoftwareEngineering #CareerGrowth #WebDevelopment #AIEngineering

How to Fix Redis Pub/Sub Message Loss

Problem

You've implemented Redis Pub/Sub in your application, expecting reliable message delivery. But you're noticing messages disappear—subscribers miss events, background jobs don't trigger, and your system behaves unpredictably. The worst part? Redis Pub/Sub has no persistence layer, so messages sent to subscribers that aren't listening are simply lost forever.

This is the fundamental challenge with Redis Pub/Sub: it's a fire-and-forget messaging system. If no one is listening when a message is published, that message vanishes into the void.


Cause

Understanding why messages are lost helps you choose the right solution:

1. Subscribers Not Connected

The most common cause. If a subscriber disconnects (network hiccup, crash, restart) and a message is published during that window, the subscriber never receives it. Redis doesn't queue messages for offline subscribers.

2. Race Conditions

Subscribers might connect after a message is published but before they subscribe to the channel. That message is already gone.

3. No Persistence

Redis Pub/Sub messages exist only in memory. If Redis crashes, all in-flight messages are lost. There's no write-ahead log or durability guarantee.

4. Slow Subscribers

If a subscriber can't keep up with the publishing rate, the Redis server's output buffer fills up. Redis may disconnect the slow subscriber, causing message loss.

5. Network Partitions

In distributed systems, network splits can cause subscribers to disconnect unexpectedly, missing messages published during the partition.


Solution

Redis Streams provide persistent, ordered message queues with consumer groups. This is the modern replacement for Pub/Sub when reliability matters.

import redis
import json
from datetime import datetime

r = redis.Redis(host='localhost', port=6379, decode_responses=True)

# Producer: Add messages to a stream
def publish_event(event_type, data):
    message = {
        'type': event_type,
        'data': json.dumps(data),
        'timestamp': datetime.utcnow().isoformat()
    }
    stream_id = r.xadd('events', message)
    print(f"Published message: {stream_id}")
    return stream_id

# Consumer: Read from stream with consumer group
def consume_events(consumer_group, consumer_name):
    # Create consumer group (idempotent)
    try:
        r.xgroup_create('events', consumer_group, id='0', mkstream=True)
    except redis.ResponseError:
        pass  # Group already exists

    while True:
        # Read pending messages first (for recovery)
        pending = r.xreadgroup(
            consumer_group,
            consumer_name,
            {'events': '>'},
            count=10,
            block=1000
        )

        if pending:
            for stream_key, messages in pending:
                for msg_id, msg_data in messages:
                    try:
                        print(f"Processing: {msg_data}")
                        # Your business logic here
                        process_message(msg_data)
                        # Acknowledge after successful processing
                        r.xack('events', consumer_group, msg_id)
                    except Exception as e:
                        print(f"Error processing {msg_id}: {e}")
                        # Message stays in pending list for retry

def process_message(msg_data):
    """Your application logic"""
    event_type = msg_data.get('type')
    data = json.loads(msg_data.get('data', '{}'))
    print(f"Event: {event_type}, Data: {data}")

# Usage
if __name__ == '__main__':
    # Publish
    publish_event('user.signup', {'user_id': 123, 'email': 'user@example.com'})

    # Consume
    consume_events('my_consumer_group', 'consumer_1')

Why Streams Win:

  • ✅ Messages persist until explicitly deleted
  • ✅ Consumer groups track which messages each consumer has processed
  • ✅ Automatic retry for failed messages
  • ✅ Scales to millions of messages
  • ✅ Supports multiple consumers processing the same stream

Option 2: Hybrid Approach (Pub/Sub + Persistence)

If you're already invested in Pub/Sub, add a fallback persistence layer:

import redis
import json
from datetime import datetime, timedelta

r = redis.Redis(host='localhost', port=6379, decode_responses=True)

class ReliablePubSub:
    def __init__(self, backup_ttl=3600):
        self.backup_ttl = backup_ttl  # Keep backups for 1 hour

    def publish(self, channel, message):
        """Publish to Pub/Sub AND backup to Redis"""
        msg_obj = {
            'data': message,
            'timestamp': datetime.utcnow().isoformat(),
            'channel': channel
        }
        msg_json = json.dumps(msg_obj)

        # Publish to subscribers
        r.publish(channel, msg_json)

        # Backup to a sorted set (for recovery)
        backup_key = f"backup:{channel}"
        score = datetime.utcnow().timestamp()
        r.zadd(backup_key, {msg_json: score})

        # Set expiration
        r.expire(backup_key, self.backup_ttl)

        return msg_json

    def subscribe_with_recovery(self, channel, callback):
        """Subscribe and replay missed messages"""
        pubsub = r.pubsub()

        # Recover missed messages from backup
        backup_key = f"backup:{channel}"
        missed_messages = r.zrange(backup_key, 0, -1)

        print(f"Recovering {len(missed_messages)} missed messages...")
        for msg_json in missed_messages:
            msg_obj = json.loads(msg_json)
            callback(msg_obj)

        # Clear backup after recovery
        r.delete(backup_key)

        # Subscribe to new messages
        pubsub.subscribe(channel)
        print(f"Subscribed to {channel}")

        for message in pubsub.listen():
            if message['type'] == 'message':
                msg_obj = json.loads(message['data'])
                callback(msg_obj)

# Usage
reliable_pub = ReliablePubSub(backup_ttl=3600)

def handle_message(msg_obj):
    print(f"Received: {msg_obj}")

# Publish
reliable_pub.publish('orders', json.dumps({'order_id': 456, 'amount': 99.99}))

# Subscribe with recovery
reliable_pub.subscribe_with_recovery('orders', handle_message)

Option 3: External Message Queue

For critical systems, use a dedicated message broker:

# Using RabbitMQ as an example
import pika
import json

def publish_with_rabbitmq(message):
    connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
    channel = connection.channel()

    # Declare durable queue
    channel.queue_declare(queue='events', durable=True)

    # Publish with persistence
    channel.basic_publish(
        exchange='',
        routing_key='events',
        body=json.dumps(message),
        properties=pika.BasicProperties(
            delivery_mode=pika.spec.PERSISTENT_DELIVERY_MODE
        )
    )
    connection.close()

def consume_with_rabbitmq():
    connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
    channel = connection.channel()

    channel.queue_declare(queue='events', durable=True)
    channel.basic_qos(prefetch_count=1)  # Process one at a time

    def callback(ch, method, properties, body):
        try:
            message = json.loads(body)
            print(f"Processing: {message}")
            # Your logic here
            ch.basic_ack(delivery_tag=method.delivery_tag)
        except Exception as e:
            print(f"Error: {e}")
            ch.basic_nack(delivery_tag=method.delivery_tag, requeue=True)

    channel.basic_consume(queue='events', on_message_callback=callback)
    channel.start_consuming()

# Usage
publish_with_rabbitmq({'event': 'user.created', 'user_id': 789})
consume_with_rabbitmq()

Tips

1. Monitor Consumer Lag

# Check how far behind consumers are
def check_consumer_lag(stream, consumer_group):
    info = r.xinfo_groups(stream)
    for group in info:
        if group['name'] == consumer_group:
            pending = group['pending']
            print(f"Pending messages: {pending}")

2. Set Client Output Buffer Limits

# In redis.conf
client-output-buffer-limit pubsub 32mb 8mb 60

Prevents slow subscribers from being disconnected.

3. Implement Idempotent Processing

def process_message_idempotent(msg_id, msg_data):
    # Check if already processed
    if r.exists(f"processed:{msg_id}"):
        return

    # Process
    do_work(msg_data)

    # Mark as processed
    r.setex(f"processed:{msg_id}", 86400, "1")

4. Use Dead Letter Queues

def handle_with_dlq(msg_id, msg_data, max_retries=3):
    retry_count = int(r.get(f"retries:{msg_id}") or 0)

    try:
        process_message(msg_data)
    except Exception as e:
        if retry_count < max_retries:
            r.incr(f"retries:{msg_id}")
        else:
            # Move to dead letter queue
            r.lpush('dlq', json.dumps(msg_data))
            print(f"Message {msg_id} moved to DLQ")

5. Health Checks

def health_check():
    try:
        r.ping()
        return True
    except:
        return False

Takeaway

Redis Pub/Sub is not a reliable message queue. It's designed for real-time notifications where occasional message loss is acceptable (e.g., live chat, notifications).

For reliable message delivery:

  • 🏆 Use Redis Streams if you want to stay within Redis
  • 🔄 Use Hybrid Pub/Sub + Backup for gradual migration
  • 🎯 Use RabbitMQ/Kafka for mission-critical systems

The choice depends on your reliability requirements, scale, and operational complexity tolerance. Start with Streams—they're battle-tested, built into Redis, and solve 95% of message queue problems without external dependencies.