Available for day contractsFrom 21st September I have availability for day and half day contracts. Please contact for more information.

Contact →
mikepreston.org

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.

Python ClientRedis ServerStringsListsSetsHashesStreamsPub/SubConsumer GroupsChannelsPython ClientRedis ServerStringsListsSetsHashesStreamsPub/SubConsumer GroupsChannels

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

Consumer 2Consumer 1Redis StreamProducerConsumer 2Consumer 1Redis StreamProducerBoth consumers receivesame messagesXADD events * dataXREAD (blocking)XREAD (blocking)Consumer 2Consumer 1Redis StreamProducerConsumer 2Consumer 1Redis StreamProducerBoth consumers receivesame messagesXADD events * dataXREAD (blocking)XREAD (blocking)

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

XACKXACKXACKStream: ordersGroup: processorsConsumer: worker-1Consumer: worker-2Consumer: worker-3XACKXACKXACKStream: ordersGroup: processorsConsumer: worker-1Consumer: worker-2Consumer: worker-3

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

PUBLISHPUBLISHPublisher 1Channel:notificationsPublisher 2Subscriber 1Subscriber 2Subscriber 3PUBLISHPUBLISHPublisher 1Channel:notificationsPublisher 2Subscriber 1Subscriber 2Subscriber 3

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