Python Redis
High-performance in-memory data store client with support for streams, pub/sub, and complex data structures.
Python Redis Cheatsheet
High-performance in-memory data store client with support for streams, pub/sub, and complex data structures.
Overview
Redis-py is the standard Python client for Redis, providing both synchronous and asynchronous interfaces. It supports all Redis data types including strings, lists, sets, sorted sets, hashes, and streams.
graph LR
A[Python Client] --> B[Redis Server]
B --> C[Strings]
B --> D[Lists]
B --> E[Sets]
B --> F[Hashes]
B --> G[Streams]
B --> H[Pub/Sub]
G --> I[Consumer Groups]
H --> J[Channels]
Connection and Basic Operations
Key Concepts
- Connection Pool: Reusable connections for better performance
- Decode Responses: Automatically decode bytes to strings
- SSL/TLS: Secure connections for production environments
- Cluster Mode: Distributed Redis across multiple nodes
Common Patterns
import redis
from redis import ConnectionPool
# Basic connection
r = redis.Redis(host='localhost', port=6379, db=0)
# With connection pool (recommended for production)
pool = ConnectionPool(
host='localhost',
port=6379,
db=0,
max_connections=10,
decode_responses=True # Return strings instead of bytes
)
r = redis.Redis(connection_pool=pool)
# Connection with authentication
r = redis.Redis(
host='redis.example.com',
port=6379,
password='secret_password',
ssl=True,
ssl_cert_reqs='required'
)
# Async connection
import redis.asyncio as aioredis
async def get_async_client():
return await aioredis.from_url(
'redis://localhost:6379',
decode_responses=True
)
Examples
# String operations
r.set('user:1:name', 'Alice')
r.set('user:1:email', 'alice@example.com', ex=3600) # Expires in 1 hour
name = r.get('user:1:name') # Returns 'Alice'
# Set with conditions
r.set('lock:resource', 'locked', nx=True) # Only if not exists
r.set('counter', 100, xx=True) # Only if exists
# Increment/Decrement
r.set('visits', 0)
r.incr('visits') # 1
r.incrby('visits', 10) # 11
r.decr('visits') # 10
# Multiple keys
r.mset({'key1': 'value1', 'key2': 'value2'})
values = r.mget(['key1', 'key2']) # ['value1', 'value2']
# Delete operations
r.delete('key1')
r.delete('key1', 'key2', 'key3') # Multiple keys
# Key existence and type
r.exists('user:1:name') # Returns 1 if exists
r.type('user:1:name') # Returns 'string'
# Pattern-based key retrieval
keys = r.keys('user:*') # Warning: blocking on large datasets
# Use SCAN for production
for key in r.scan_iter('user:*', count=100):
print(key)
Stream Operations
Key Concepts
- Streams: Append-only log data structure for event sourcing
- Entry ID: Unique identifier (timestamp-sequence) for each entry
- Blocking Reads: Efficient waiting for new messages
- Trimming: Limit stream size by count or ID
Common Patterns
sequenceDiagram
participant P as Producer
participant S as Redis Stream
participant C1 as Consumer 1
participant C2 as Consumer 2
P->>S: XADD events * data
S-->>C1: XREAD (blocking)
S-->>C2: XREAD (blocking)
Note over S: Both consumers receive<br/>same messages
Examples
# Add entries to stream
# Auto-generate ID with '*'
entry_id = r.xadd('events', {'type': 'click', 'user_id': '123'})
# Returns: '1699012345678-0'
# Specify maximum length (approximate)
r.xadd('events', {'data': 'payload'}, maxlen=1000)
# Exact maximum length
r.xadd('events', {'data': 'payload'}, maxlen=1000, approximate=False)
# XREAD - Read entries from streams
# Read all entries from beginning
entries = r.xread({'events': '0-0'})
# Returns: [['events', [('1699012345678-0', {'type': 'click', 'user_id': '123'})]]]
# Read only new entries (blocking)
entries = r.xread({'events': '$'}, block=5000) # Block for 5 seconds
# Read from multiple streams
entries = r.xread({
'stream1': '0-0',
'stream2': '0-0'
}, count=10)
# XRANGE - Read entries by ID range
entries = r.xrange('events', min='-', max='+') # All entries
entries = r.xrange('events', min='1699012345678-0', max='+', count=100)
# XREVRANGE - Reverse order
entries = r.xrevrange('events', max='+', min='-', count=10)
# XLEN - Get stream length
length = r.xlen('events')
# XTRIM - Trim stream
r.xtrim('events', maxlen=1000)
r.xtrim('events', minid='1699012345678-0') # Remove entries before ID
# XDEL - Delete specific entries
r.xdel('events', '1699012345678-0')
# XINFO - Stream information
info = r.xinfo_stream('events')
# Returns: {'length': 100, 'first-entry': ..., 'last-entry': ...}
Consumer Groups
Key Concepts
- Consumer Group: Named group that tracks message delivery
- Pending Entries List (PEL): Messages delivered but not acknowledged
- Message Claiming: Reassign stuck messages to other consumers
- Acknowledgement: Confirm message processing completion
Common Patterns
graph TD
S[Stream: orders] --> G1[Group: processors]
G1 --> C1[Consumer: worker-1]
G1 --> C2[Consumer: worker-2]
G1 --> C3[Consumer: worker-3]
C1 --> |XACK| G1
C2 --> |XACK| G1
C3 --> |XACK| G1
style S fill:#e1f5fe
style G1 fill:#fff3e0
Examples
# Create consumer group
# Start from beginning of stream
try:
r.xgroup_create('orders', 'processors', id='0', mkstream=True)
except redis.ResponseError as e:
if 'BUSYGROUP' not in str(e):
raise
# Start from new messages only
r.xgroup_create('orders', 'processors', id='$', mkstream=True)
# Read messages as consumer
# '>' means only new messages not yet delivered
messages = r.xreadgroup(
groupname='processors',
consumername='worker-1',
streams={'orders': '>'},
count=10,
block=5000
)
# Process and acknowledge messages
for stream, entries in messages:
for entry_id, data in entries:
try:
process_order(data)
r.xack('orders', 'processors', entry_id)
except Exception as e:
# Message remains in pending list
log_error(e)
# Check pending messages
pending = r.xpending('orders', 'processors')
# Returns: {'pending': 5, 'min': '...', 'max': '...', 'consumers': [...]}
# Detailed pending info
pending_detail = r.xpending_range(
'orders', 'processors',
min='-', max='+',
count=10,
consumername='worker-1'
)
# Claim stuck messages (idle for 60 seconds)
claimed = r.xclaim(
'orders', 'processors', 'worker-2',
min_idle_time=60000, # milliseconds
message_ids=['1699012345678-0']
)
# Auto-claim with XAUTOCLAIM
# Returns [next_start_id, claimed_entries, deleted_ids]
next_id, claimed, deleted_ids = r.xautoclaim(
'orders', 'processors', 'worker-2',
min_idle_time=60000,
start_id='0-0',
count=10
)
# Delete consumer from group
r.xgroup_delconsumer('orders', 'processors', 'worker-1')
# Delete consumer group
r.xgroup_destroy('orders', 'processors')
# Consumer group info
groups = r.xinfo_groups('orders')
consumers = r.xinfo_consumers('orders', 'processors')
Pub/Sub Patterns
Key Concepts
- Channel: Named message destination
- Pattern Subscription: Subscribe using wildcards
- Message Handler: Callback for incoming messages
- Non-blocking: Messages not persisted; missed if not subscribed
Common Patterns
graph LR
P1[Publisher 1] --> |PUBLISH| C[Channel: notifications]
P2[Publisher 2] --> |PUBLISH| C
C --> S1[Subscriber 1]
C --> S2[Subscriber 2]
C --> S3[Subscriber 3]
style C fill:#e8f5e9
Examples
# Publisher
def publish_notification(channel, message):
r.publish(channel, message)
publish_notification('notifications', 'New order received')
publish_notification('alerts:critical', 'Server down!')
# Subscriber - synchronous
pubsub = r.pubsub()
# Subscribe to specific channels
pubsub.subscribe('notifications', 'alerts')
# Subscribe with pattern
pubsub.psubscribe('alerts:*')
# Listen for messages (blocking)
for message in pubsub.listen():
if message['type'] == 'message':
print(f"Channel: {message['channel']}")
print(f"Data: {message['data']}")
elif message['type'] == 'pmessage':
print(f"Pattern: {message['pattern']}")
print(f"Channel: {message['channel']}")
print(f"Data: {message['data']}")
# Non-blocking message retrieval
message = pubsub.get_message()
if message and message['type'] == 'message':
process_message(message['data'])
# Handler-based subscription
def notification_handler(message):
print(f"Received: {message['data']}")
def alert_handler(message):
send_alert(message['data'])
pubsub.subscribe(**{
'notifications': notification_handler,
'alerts': alert_handler
})
# Run in thread
thread = pubsub.run_in_thread(sleep_time=0.001)
# ... do other work ...
thread.stop()
# Async pub/sub
import redis.asyncio as aioredis
async def subscriber():
r = await aioredis.from_url('redis://localhost')
pubsub = r.pubsub()
await pubsub.subscribe('notifications')
async for message in pubsub.listen():
if message['type'] == 'message':
await process_async(message['data'])
# Unsubscribe
pubsub.unsubscribe('notifications')
pubsub.punsubscribe('alerts:*')
pubsub.close()
Data Structures
Key Concepts
- Lists: Ordered collections, ideal for queues
- Sets: Unordered unique elements
- Sorted Sets: Scored elements for rankings/leaderboards
- Hashes: Field-value maps for objects
Examples
Lists
# Push operations
r.lpush('queue', 'task1', 'task2') # Left push (prepend)
r.rpush('queue', 'task3') # Right push (append)
# Pop operations
task = r.lpop('queue') # Remove from left
task = r.rpop('queue') # Remove from right
task = r.blpop('queue', timeout=5) # Blocking pop
# Range operations
items = r.lrange('queue', 0, -1) # All items
items = r.lrange('queue', 0, 9) # First 10
# Queue pattern (FIFO)
r.rpush('jobs', 'job1')
job = r.lpop('jobs')
# Stack pattern (LIFO)
r.lpush('stack', 'item1')
item = r.lpop('stack')
# List length
length = r.llen('queue')
# Set by index
r.lset('queue', 0, 'updated_task')
# Trim list
r.ltrim('queue', 0, 99) # Keep first 100 elements
Sets
# Add members
r.sadd('tags', 'python', 'redis', 'database')
# Remove members
r.srem('tags', 'database')
# Check membership
r.sismember('tags', 'python') # True
# Get all members
members = r.smembers('tags')
# Set operations
r.sadd('user:1:skills', 'python', 'javascript')
r.sadd('user:2:skills', 'python', 'go')
# Intersection
common = r.sinter('user:1:skills', 'user:2:skills') # {'python'}
# Union
all_skills = r.sunion('user:1:skills', 'user:2:skills')
# Difference
unique = r.sdiff('user:1:skills', 'user:2:skills') # {'javascript'}
# Random member
random_tag = r.srandmember('tags')
# Pop random member
popped = r.spop('tags')
# Cardinality
count = r.scard('tags')
Sorted Sets
# Add with scores
r.zadd('leaderboard', {'player1': 100, 'player2': 85, 'player3': 92})
# Update score
r.zadd('leaderboard', {'player1': 105})
# Increment score
r.zincrby('leaderboard', 10, 'player2') # Now 95
# Get rank (0-based, ascending)
rank = r.zrank('leaderboard', 'player1')
# Get rank (descending)
rank = r.zrevrank('leaderboard', 'player1')
# Get score
score = r.zscore('leaderboard', 'player1')
# Range by rank (ascending)
top_players = r.zrange('leaderboard', 0, 9, withscores=True)
# Range by rank (descending)
top_players = r.zrevrange('leaderboard', 0, 9, withscores=True)
# Range by score
players = r.zrangebyscore('leaderboard', 90, 100, withscores=True)
# Remove members
r.zrem('leaderboard', 'player3')
# Remove by rank
r.zremrangebyrank('leaderboard', 0, 4) # Remove bottom 5
# Remove by score
r.zremrangebyscore('leaderboard', 0, 50) # Remove scores 0-50
# Cardinality
count = r.zcard('leaderboard')
# Count by score range
count = r.zcount('leaderboard', 90, 100)
Hashes
# Set fields
r.hset('user:100', 'name', 'Alice')
r.hset('user:100', mapping={
'email': 'alice@example.com',
'age': 30,
'city': 'London'
})
# Get field
name = r.hget('user:100', 'name')
# Get multiple fields
values = r.hmget('user:100', ['name', 'email'])
# Get all fields and values
user = r.hgetall('user:100')
# {'name': 'Alice', 'email': 'alice@example.com', ...}
# Check field exists
r.hexists('user:100', 'name') # True
# Delete fields
r.hdel('user:100', 'age')
# Get all field names
fields = r.hkeys('user:100')
# Get all values
values = r.hvals('user:100')
# Increment numeric field
r.hincrby('user:100', 'login_count', 1)
r.hincrbyfloat('user:100', 'balance', 10.50)
# Field count
count = r.hlen('user:100')
# Set if not exists
r.hsetnx('user:100', 'created_at', '2024-01-01')
Expiration and TTL
Key Concepts
- TTL: Time-to-live in seconds
- PTTL: Precision TTL in milliseconds
- EXPIREAT: Expire at Unix timestamp
- PERSIST: Remove expiration
Examples
# Set expiration on new key
r.set('session:abc123', 'data', ex=3600) # Expires in 1 hour
r.set('session:abc123', 'data', px=60000) # Expires in 60 seconds (ms)
# Set expiration on existing key
r.expire('session:abc123', 3600) # 1 hour
r.pexpire('session:abc123', 60000) # 60 seconds (ms)
# Set expiration at specific time
import time
r.expireat('session:abc123', int(time.time()) + 3600)
# Get remaining TTL
ttl = r.ttl('session:abc123') # Seconds (-1 if no expiry, -2 if not exists)
pttl = r.pttl('session:abc123') # Milliseconds
# Remove expiration
r.persist('session:abc123')
# Check if key has expiration
ttl = r.ttl('session:abc123')
has_expiry = ttl > 0
# Set with conditional expiration
r.set('lock', 'value', ex=30, nx=True) # Lock with 30s timeout
# Update expiration only
r.expire('cache:data', 3600, xx=True) # Only if key has expiry
r.expire('cache:data', 3600, nx=True) # Only if key has no expiry
r.expire('cache:data', 3600, gt=True) # Only if new TTL > current
r.expire('cache:data', 3600, lt=True) # Only if new TTL < current
Pipeline and Transactions
Key Concepts
- Pipeline: Batch commands for network efficiency
- Transaction: Atomic execution with MULTI/EXEC
- WATCH: Optimistic locking for check-and-set
- Lua Scripts: Server-side atomic operations
Examples
Pipeline (Non-atomic batching)
# Basic pipeline
pipe = r.pipeline(transaction=False)
pipe.set('key1', 'value1')
pipe.set('key2', 'value2')
pipe.get('key1')
pipe.get('key2')
results = pipe.execute()
# ['OK', 'OK', 'value1', 'value2']
# Context manager
with r.pipeline(transaction=False) as pipe:
for i in range(1000):
pipe.set(f'key:{i}', f'value:{i}')
results = pipe.execute()
# Chained pipeline
results = (
r.pipeline(transaction=False)
.set('a', 1)
.set('b', 2)
.incr('a')
.incr('b')
.execute()
)
Transaction (Atomic)
# Basic transaction
# Note: transaction=True (the default) already wraps execute() in MULTI/EXEC
# automatically — you do NOT need to call pipe.multi() in this case.
# pipe.multi() switches to immediate-MULTI mode and is only needed after
# pipe.watch() calls (optimistic locking), as shown in the WATCH example below.
pipe = r.pipeline(transaction=True) # Default
pipe.set('key1', 'value1')
pipe.incr('counter')
pipe.execute() # Executes atomically
# WATCH for optimistic locking
def transfer_funds(from_account, to_account, amount):
with r.pipeline() as pipe:
while True:
try:
# Watch keys for changes
pipe.watch(from_account, to_account)
# Get current balances
from_balance = int(r.get(from_account) or 0)
to_balance = int(r.get(to_account) or 0)
if from_balance < amount:
pipe.unwatch()
raise ValueError('Insufficient funds')
# Start transaction
pipe.multi()
pipe.set(from_account, from_balance - amount)
pipe.set(to_account, to_balance + amount)
pipe.execute()
break
except redis.WatchError:
# Key was modified, retry
continue
# Discard transaction
pipe = r.pipeline()
pipe.multi()
pipe.set('key', 'value')
pipe.discard() # Cancel transaction
Lua Scripts
# Register script
increment_script = r.register_script("""
local current = redis.call('GET', KEYS[1])
if current then
current = tonumber(current)
else
current = 0
end
local new_value = current + tonumber(ARGV[1])
redis.call('SET', KEYS[1], new_value)
return new_value
""")
# Execute script
result = increment_script(keys=['counter'], args=[10])
# Rate limiter script
rate_limit_script = r.register_script("""
local key = KEYS[1]
local limit = tonumber(ARGV[1])
local window = tonumber(ARGV[2])
local current = redis.call('INCR', key)
if current == 1 then
redis.call('EXPIRE', key, window)
end
if current > limit then
return 0
end
return 1
""")
def check_rate_limit(user_id, limit=100, window=60):
key = f'rate_limit:{user_id}'
allowed = rate_limit_script(keys=[key], args=[limit, window])
return bool(allowed)
Quick Reference
| Operation | Command | Example |
|---|---|---|
| Connection | ||
| Connect | redis.Redis() |
r = redis.Redis(host='localhost', decode_responses=True) |
| Pool | ConnectionPool() |
pool = ConnectionPool(max_connections=10) |
| Strings | ||
| Set | set() |
r.set('key', 'value', ex=3600) |
| Get | get() |
r.get('key') |
| Delete | delete() |
r.delete('key1', 'key2') |
| Increment | incr() |
r.incr('counter') |
| Streams | ||
| Add | xadd() |
r.xadd('stream', {'field': 'value'}) |
| Read | xread() |
r.xread({'stream': '$'}, block=5000) |
| Group read | xreadgroup() |
r.xreadgroup('group', 'consumer', {'stream': '>'}) |
| Acknowledge | xack() |
r.xack('stream', 'group', entry_id) |
| Consumer Groups | ||
| Create | xgroup_create() |
r.xgroup_create('stream', 'group', '0', mkstream=True) |
| Claim | xclaim() |
r.xclaim('stream', 'group', 'consumer', 60000, [id]) |
| Pending | xpending() |
r.xpending('stream', 'group') |
| Pub/Sub | ||
| Publish | publish() |
r.publish('channel', 'message') |
| Subscribe | subscribe() |
pubsub.subscribe('channel') |
| Pattern sub | psubscribe() |
pubsub.psubscribe('channel:*') |
| Lists | ||
| Push | lpush()/rpush() |
r.rpush('list', 'item') |
| Pop | lpop()/rpop() |
r.lpop('list') |
| Range | lrange() |
r.lrange('list', 0, -1) |
| Sets | ||
| Add | sadd() |
r.sadd('set', 'member') |
| Members | smembers() |
r.smembers('set') |
| Intersect | sinter() |
r.sinter('set1', 'set2') |
| Sorted Sets | ||
| Add | zadd() |
r.zadd('zset', {'member': 100}) |
| Range | zrange() |
r.zrange('zset', 0, 9, withscores=True) |
| Rank | zrank() |
r.zrank('zset', 'member') |
| Hashes | ||
| Set | hset() |
r.hset('hash', mapping={'k': 'v'}) |
| Get | hget()/hgetall() |
r.hgetall('hash') |
| Increment | hincrby() |
r.hincrby('hash', 'field', 1) |
| TTL | ||
| Set expiry | expire() |
r.expire('key', 3600) |
| Get TTL | ttl() |
r.ttl('key') |
| Remove | persist() |
r.persist('key') |
| Pipeline | ||
| Create | pipeline() |
pipe = r.pipeline() |
| Execute | execute() |
results = pipe.execute() |
| Watch | watch() |
pipe.watch('key') |
Common Issues and Solutions
Connection Issues
| Issue | Cause | Solution |
|---|---|---|
ConnectionError: Error connecting to localhost:6379 |
Redis not running | Start Redis: redis-server or systemctl start redis |
AuthenticationError |
Wrong password | Verify password parameter matches Redis config |
TimeoutError |
Network issues or slow Redis | Increase socket_timeout parameter |
| Connection pool exhausted | Too many concurrent connections | Increase max_connections or use connection pooling |
Data Type Errors
# Issue: ResponseError: WRONGTYPE Operation against a key holding the wrong kind of value
# Solution: Check key type before operations
if r.type('key') == 'string':
r.get('key')
elif r.type('key') == 'list':
r.lrange('key', 0, -1)
Memory Issues
# Issue: OOM command not allowed when used memory > 'maxmemory'
# Solution: Set eviction policy and monitor memory
# Check memory usage
info = r.info('memory')
print(f"Used memory: {info['used_memory_human']}")
# Set maxmemory policy (in redis.conf or via CONFIG)
r.config_set('maxmemory-policy', 'allkeys-lru')
Stream Consumer Group Issues
# Issue: BUSYGROUP Consumer Group name already exists
# Solution: Check before creating
try:
r.xgroup_create('stream', 'group', '0', mkstream=True)
except redis.ResponseError as e:
if 'BUSYGROUP' in str(e):
pass # Group already exists
else:
raise
# Issue: NOGROUP No such consumer group
# Solution: Create group with mkstream=True
r.xgroup_create('stream', 'group', '0', mkstream=True)
Pipeline/Transaction Errors
# Issue: WatchError during transaction
# Solution: Implement retry logic
def safe_transaction(max_retries=3):
for attempt in range(max_retries):
try:
with r.pipeline() as pipe:
pipe.watch('key')
# ... transaction logic ...
pipe.execute()
break
except redis.WatchError:
if attempt == max_retries - 1:
raise
continue
Encoding Issues
# Issue: TypeError: a bytes-like object is required
# Solution: Use decode_responses=True
r = redis.Redis(host='localhost', decode_responses=True)
# Or decode manually
value = r.get('key')
if value:
value = value.decode('utf-8')
Performance Issues
# Issue: Slow KEYS command blocking Redis
# Solution: Use SCAN instead
# Bad
keys = r.keys('pattern:*')
# Good
keys = list(r.scan_iter('pattern:*', count=100))
# Issue: Too many round trips
# Solution: Use pipeline
with r.pipeline(transaction=False) as pipe:
for key in keys:
pipe.get(key)
values = pipe.execute()
Pub/Sub Message Loss
# Issue: Messages missed when subscriber disconnects
# Solution: Use Streams instead for durability
# Or implement reconnection logic
def resilient_subscriber():
while True:
try:
pubsub = r.pubsub()
pubsub.subscribe('channel')
for message in pubsub.listen():
process(message)
except redis.ConnectionError:
time.sleep(1)
continue