Redis PubSub Broker
The Redis PubSub broker uses Redis Lists for simple queuing with consumer groups. Despite its name, it doesn't use Redis's native Pub/Sub mechanism but implements a queue-based messaging system.
Overview
This broker uses Redis Lists as queues where:
- Publishers append messages to Redis lists using
rpush - Subscribers read messages from lists using
lindexwith position tracking - Messages are persisted in Redis lists until Redis is flushed
- Consumer groups track their position in each list independently
Installation
The Redis PubSub broker requires the redis Python package:
# Install Protean with Redis support
pip install "protean[redis]"
# Or install Redis package separately
pip install "redis>=8.0.0,<8.2.0"
Configuration
[brokers.notifications]
provider = "redis_pubsub"
URI = "redis://localhost:6379/0"
Configuration Options
| Option | Default | Description |
|---|---|---|
provider |
Required | Must be "redis_pubsub" for Redis PubSub broker |
URI |
Required | Redis connection string |
Capabilities
The Redis PubSub broker provides simple queuing capabilities:
- ✅ BASIC_PUBSUB - Publish and subscribe
- ✅ SIMPLE_QUEUING - Consumer groups with position tracking
- ❌ RELIABLE_MESSAGING - No acknowledgments (ack/nack not supported)
- ❌ ORDERED_MESSAGING - No ordering guarantees
- ❌ ENTERPRISE_STREAMING - Not supported
Usage Examples
The examples on this page need a running Redis server. They build on one another.
Basic Publishing
Configure the broker, then publish to a stream. Each stream is a Redis list with the same name:
# fragment
import os
from protean import Domain
domain = Domain(name="Notifications")
domain.config["brokers"]["notifications"] = {
"provider": "redis_pubsub",
"URI": os.environ.get("REDIS_URL", "redis://localhost:6379/0"),
}
# Call subscribers in the publishing process, as soon as a message is published
domain.config["message_processing"] = "sync"
domain.init(traverse=False)
with domain.domain_context():
# RPUSH the notification onto the "user:notifications" list
domain.brokers["notifications"].publish(
stream="user:notifications",
message={
"type": "notification",
"user_id": "123",
"title": "New Message",
"body": "You have a new message!",
},
)
Subscribing to Streams
A subscriber is a class with a __call__ method that receives each message as
a dict. Name the broker with broker=, and register the subscriber before you
initialize the domain:
# fragment
pushed = []
@domain.subscriber(stream="user:notifications", broker="notifications")
class NotificationSubscriber:
def __call__(self, payload: dict) -> None:
# Send push notification to the user's device
pushed.append((payload["user_id"], payload["title"]))
With message_processing = "sync" (set in the configuration above), the
subscriber runs as soon as the message is published. Here pushed becomes
[("123", "New Message")].
A subscriber reads one stream by its exact name. Stream names are not
patterns, so stream="chat:*" reads a list literally named chat:*.
Consumer Groups
Each consumer group keeps its own position in the list, in a Redis key named
position:<stream>:<group>. Two groups read the same messages:
# fragment
with domain.domain_context():
broker = domain.brokers["notifications"]
broker.publish("orders", {"type": "order.created", "order_id": "A1"})
# Each group keeps its own position in the list, so both read the message
_, billing_message = broker.get_next("orders", "billing")
shipping_id, shipping_message = broker.get_next("orders", "shipping")
# Acknowledgment is not supported: ack() returns False
acknowledged = broker.ack("orders", shipping_id, "shipping")
The broker does not track whether a message was processed. ack and nack
log a warning and return False.
Limitations and Considerations
Limited Persistence
Messages stay in Redis lists until Redis is flushed. If Redis restarts without persistence (RDB snapshots or AOF) turned on, the messages are lost. For durable messaging, use the Redis Streams broker with persistence turned on.
No Acknowledgment Support
publish returns a message identifier, but there is no delivery confirmation.
A consumer group moves its position forward when it reads a message, whether
or not a subscriber handled it.
Position Tracking
The group positions are Redis keys. If they are deleted (for example by
FLUSHDB), every group reads its streams from the start again.
Performance Considerations
Message Size
Keep messages small. Put large payloads in a store such as S3, and publish a reference to them (a URL and a size) instead of the payload.
Stream Naming
Use hierarchical names such as user:123:notifications or
system:alerts:critical, so redis-cli --scan --pattern can find related
streams.
Message Processing
Messages are processed sequentially per consumer group. Each group maintains its own position counter in Redis.
Monitoring and Debugging
Redis CLI Commands
The broker stores messages in lists, so the list commands show its state:
# List streams and group positions
redis-cli --scan --pattern "*"
# Count the messages in a stream
redis-cli LLEN user:notifications
# Read the first ten messages
redis-cli LRANGE user:notifications 0 9
# Show a group's position in a stream
redis-cli GET position:user:notifications:billing
Logging
The broker logs through the protean.adapters.broker.redis_pubsub logger. Set
it to DEBUG with the standard logging module when you troubleshoot.
Health Checks
health_stats() returns the common health dict. The Redis details, such as
connected_clients and used_memory_human, are under details:
# fragment
def health_check():
broker = domain.brokers["notifications"]
# Test connectivity
if not broker.ping():
return {"status": "unhealthy", "error": "Connection failed"}
# Get health statistics; the Redis details sit under "details"
stats = broker.health_stats()
details = stats["details"]
return {
"status": stats["status"],
"connected_clients": details["connected_clients"],
"used_memory": details["used_memory_human"],
}
with domain.domain_context():
health = health_check()
Migration Strategies
To Redis Streams
When you need persistence and reliability, change the provider:
# Before: Redis PubSub
[brokers.notifications]
provider = "redis_pubsub"
URI = "redis://localhost:6379/0"
# After: Redis Streams
[brokers.notifications]
provider = "redis"
URI = "redis://localhost:6379/0"
Subscribers stay the same. Two things change:
ackandnackwork, so a message that is not acknowledged stays pending for its group.publishreturns a Redis Streams entry id such as1700000000000-0.
Working with Redis PubSub
Use for Simple Queuing
Redis PubSub broker is suitable for:
- Simple message distribution
- Development and testing
- Scenarios where ack/nack isn't needed
- Basic consumer group functionality
Not suitable for:
- Critical business events requiring acknowledgment
- Complex message routing
- Scenarios requiring message replay
- High-throughput production systems
Comparison with Other Brokers
| Feature | Redis PubSub | Redis Streams | Inline |
|---|---|---|---|
| Persistence | Redis Lists | Yes (durable) | No |
| Delivery Guarantee | None | At-least-once | Best-effort |
| Consumer Groups | Basic | Advanced | Yes |
| Message Ordering | No | Yes | No |
| Acknowledgments | No | Yes | Yes |
| Performance | High | High | Very High |
| Use Case | Simple Queuing | Event Streaming | Development |
Related pages
- Learn about Redis Streams broker for reliable messaging
- Understand broker capabilities in detail
- Explore custom broker development