The Tool Desk
Outbyte Driver Updater FREEScan for outdated or missing drivers - takes under a minuteDriver Scan →Outbyte PC Repair FREEClear out junk files and repair common Windows errorsFree Scan →iTechGuides is reader-supported. When you buy through links on our site, we may earn an affiliate commission. As an Amazon Associate I earn from qualifying purchases. Learn more
For events that must survive a worker crash, be shared across a worker pool, and stay replayable for a chosen period, the pattern Redis documents for Python is a Redis Stream read through a consumer group. Producers append with XADD. Each worker reads new entries with XREADGROUP and the > ID, calls XACK only after its work succeeds, and a recovery step uses XPENDING and XAUTOCLAIM to take over deliveries that stayed unacknowledged too long. The guarantee is at-least-once processing rather than exactly-once side effects, a distinction the sections below explain.
WRedis is a separate package on PyPI that offers a decorator-based wrapper around Redis Streams. Its project page documents the methods, but it does not describe how acknowledgement, pending-entry recovery, or retention behave when a worker fails. The sections on WRedis explain what can and cannot be confirmed from that page.
What a Redis Stream guarantees
Redis describes a stream as “an append-only log of field/value entries with auto-generated, time-ordered IDs” (Redis streaming guide for redis-py). Three properties shape the design:
Crashes, No Sound, or Screen Glitches?
Random freezes, missing sound and display glitches usually trace back to one bad driver. Find and replace yours safely.Free scan · under a minuteWindows Errors? Fix Them Before They Spread
Repair common Windows errors and clear accumulated junk for a smoother, more stable PC - no reinstall needed.Free scan · no reinstall- Entries are not edited in place. They leave a stream only through trimming or explicit deletion.
- IDs carry a millisecond timestamp before the dash, so an ID cutoff doubles as a time cutoff. This matters when you set retention.
- Reading a range changes no group’s position.
XRANGEreturns entries without advancing any consumer group, which is how you replay history independently of processing.
The worker lifecycle, step by step
A reliable worker moves each event through five stages. Appending and creating the group happen once; reading, acknowledging, and recovering run in a loop.
#1 Best Overall
1. Append events with XADD
Producers add structured events with XADD. Include an idempotency key in each payload so consumers can recognise a repeat later.
import redis
r = redis.Redis(host='localhost', port=6379, decode_responses=True)
entry_id = r.xadd(
'orders:events',
{'type': 'order_placed', 'order_id': '1042', 'idempotency_key': 'order-1042-placed'},
)
print(entry_id) # for example: 1765432100000-0
The return value is the generated entry ID. Store it if downstream systems need to refer to a specific event.
2. Create the consumer group once
Create the group with XGROUP CREATE, adding MKSTREAM so the stream is created if it does not exist yet. If the group already exists, Redis returns a BUSYGROUP error. The starting ID matters only at creation; the next-but-one section covers the choice.
What’s actually slowing this PC down?
Pick the symptom - the matching free tool is one click away.
from redis.exceptions import ResponseError
try:
r.xgroup_create('orders:events', 'billing', id='0', mkstream=True)
except ResponseError as exc:
if 'BUSYGROUP' not in str(exc):
raise # the group already exists; any other error is real
3. Read as a group member with XREADGROUP
The > ID asks for entries never delivered to this group. Each entry goes to one member of the group, which is how workers share the load. Every delivery is recorded in the group’s pending entries list (PEL) under the name of the consumer that received it.
Reading with an explicit ID such as 0 instead returns that consumer’s own pending entries, which is useful when a worker restarts. The block argument waits for new entries without a busy loop; when the timeout passes, no batch is returned.
batches = r.xreadgroup(
'billing', 'worker-1', {'orders:events': '>'}, count=10, block=5000
)
4. Acknowledge only after the work succeeds
The acknowledgement boundary is the core of the design. XACK removes an entry from the PEL, and it must run after the side effect has completed.
Rank #2
def handle(fields):
# Must be safe to run twice for the same idempotency_key.
...
def run_worker(consumer_name):
while True:
batches = r.xreadgroup(
'billing', consumer_name, {'orders:events': '>'}, count=10, block=5000
)
for stream, messages in batches or []:
for msg_id, fields in messages:
handle(fields) # an exception skips the next line
r.xack('orders:events', 'billing', msg_id)
If handle raises, XACK never runs and the entry stays pending. If the process dies after the side effect but before XACK, the entry is redelivered later. That window is what idempotent handlers have to cover.
Quick wins for a faster PC:
Fix the driver behind crashes, sound loss and screen glitchesFind Drivers →Clear out junk files and repair common Windows errorsFree Scan →Scan for outdated or missing drivers - takes under a minuteDriver Scan →5. Recover idle deliveries with XPENDING and XAUTOCLAIM
A recovery worker, or a periodic task inside each worker, inspects the PEL and claims entries that have been idle long enough. Check the summary first, then the per-entry view.
# Summary: pending count, lowest and highest pending ID, counts per consumer
print(r.xpending('orders:events', 'billing'))
# Per entry: current owner, time since delivery in ms, delivery count
for entry in r.xpending_range('orders:events', 'billing', min='-', max='+', count=20):
print(entry['message_id'], entry['consumer'],
entry['time_since_delivered'], entry['times_delivered'])
# Claim entries idle for at least 120 seconds
cursor, claimed, deleted_ids = r.xautoclaim(
'orders:events', 'billing', 'worker-2',
min_idle_time=120_000, start_id='0-0', count=50,
)
XAUTOCLAIM returns a cursor for the next scan, the claimed entries, and the IDs of pending entries that have since been removed from the stream. Keep calling it with the returned cursor until the cursor comes back as 0-0. Claimed entries then go through the same handle-then-acknowledge path as new ones, under the new consumer’s name.
Processing guarantees and idempotent handlers
Acknowledging after the work gives at-least-once behaviour. The outcomes depend on where a crash happens:
- Crash before the handler finishes: the entry stays pending and is reclaimed later.
- Crash after the external write but before
XACK: the entry is reclaimed and the write happens again. - Crash after
XACK: nothing is redelivered.
Make the second case harmless. Record the event’s idempotency key in the same transaction as the business update, and skip any event whose key is already present. Where the operation is naturally idempotent, such as setting a status to a fixed value, a separate key is unnecessary.
Consumer groups or direct readers
Plain XREAD suits tailing a stream, and XREADGROUP suits shared, recoverable processing. The difference is in the bookkeeping:
| Question | XREAD | XREADGROUP |
|---|---|---|
| Who receives each entry | Every reader gets the entries it asks for, tracking its own last ID | Each entry goes to one member of the group |
| Pending-entry tracking | None | Recorded per consumer in the PEL |
| Acknowledgement | Not available | XACK removes entries from pending state |
| Recovery of failed deliveries | Not provided by Redis | XAUTOCLAIM transfers idle entries |
| Typical use | Dashboards, debugging, direct tailing | Shared work queues with recovery |
Separate groups give independent passes over the same events. A billing group and an analytics group can each read every order event, and each keeps its own cursor and PEL. Members of one group share work; groups do not share it with each other.
Starting position and replay
The starting ID decides what a new group processes. It is fixed when the group is created:
| Starting ID | What the group processes | Use when |
|---|---|---|
0 |
Every retained entry, then new arrivals | The group must backfill from retained history |
$ |
Only entries appended after creation | Only future events matter |
| A specific entry ID | Entries from that ID onward | You have a known cutoff, such as a deployment time |
For replay without touching any group, read a range directly. This does not advance a cursor or create pending state:
Do these 3 things before closing this tab:
1Scan for outdated or missing drivers - takes under a minute2Clear out junk files and repair common Windows errors3Fix the driver behind crashes, sound loss and screen glitches# The full retained history, up to 500 entries at a time
history = r.xrange('orders:events', min='-', max='+', count=500)
# Entries from a known millisecond timestamp onward
window = r.xrange('orders:events', min='1765400000000-0', max='+', count=500)
Replay only reaches what retention still holds, which is the subject of the next section.
Retention: bounding history without cutting the replay window
Trimming is how a stream’s size is bounded. The two methods bound different things:
| Method | Option | Bounds by | Behaviour to plan for |
|---|---|---|---|
| Approximate length cap | MAXLEN with approximate trimming |
Entry count | Redis removes entries in groups, so the stream can hold somewhat more than the cap; the cap is not exact |
| Minimum ID | MINID |
Time, through a millisecond ID cutoff | Removes every entry below the cutoff |
import time
# Cap by approximate entry count
r.xadd('orders:events', {'type': 'order_placed'}, maxlen=100_000, approximate=True)
# Keep roughly seven days, using an ID cutoff
cutoff_ms = int((time.time() - 7 * 86400) * 1000)
r.xtrim('orders:events', minid=f'{cutoff_ms}-0')
Choose the method from the replay requirement. Trim by time when a fixed replay period matters, such as re-driving a week of billing events. Trim by count when memory is the constraint and the window can be described as a number of events. Trimming can also remove entries that are still in a pending list. Set retention longer than your longest expected recovery delay, or a stuck delivery may point to a payload that no longer exists.
Monitoring lag and pending entries
Two numbers answer different questions. Lag counts entries not yet delivered to a group, so it shows whether processing keeps up with arrivals. Pending counts entries delivered but not acknowledged, so it shows whether the handle-then-acknowledge path is working. Inspect them with:
XINFO GROUPS orders:events
XINFO CONSUMERS orders:events billing
XPENDING orders:events billing - + 20
XINFO GROUPS reports each group’s lag and pending count, XINFO CONSUMERS shows each consumer’s pending count and idle time, and the extended XPENDING form lists individual entries.
| Signal | Usually means | First check |
|---|---|---|
| Lag rising, pending low | Producers are outpacing processing | Number of active consumers and per-message handler time |
| Pending rising, workers running | Entries are delivered but not acknowledged | Whether XACK is skipped after exceptions |
| Pending concentrated on one consumer | That worker crashed or is stuck | Its idle time in XINFO CONSUMERS, then claim its entries |
| One entry with a high delivery count | A message that fails on every attempt | The payload; route it to an application-level dead-letter stream |
How long a delivery should sit idle before reclaim
The idle threshold passed to XAUTOCLAIM is the most common cause of duplicate concurrent processing. If it is shorter than a legitimate handler run, a second worker claims an entry the first is still processing, and both proceed.
- Measure handler duration under realistic load, using a high percentile such as p99, and set the threshold comfortably above it.
- Redis has no lease for a pending entry, so the idle threshold is the only control. If some handlers run for minutes, split that work into smaller events rather than lengthening the threshold for everything.
- Because the threshold is shared by all workers in the group, a single long job can force a longer threshold on the whole pool.
Scaling past one stream key
A stream is a single key, and a key lives on one shard in Redis Cluster. Adding consumers to a group spreads delivery across workers, but it does not add write capacity to the key. Partition the stream when one of these holds:
- the write rate of one stream approaches what its shard can absorb;
- ordering matters only within a tenant or entity, so a stream per tenant or per entity is enough.
Partitioning changes ordering guarantees. Events are ordered within a partition only, and any consumer that needs order across partitions must merge them itself. Give independent groups separate consumer pools when one group’s workload must not take workers from another; a shared pool is fine when the groups should share capacity.
Version and client compatibility
The version requirements below come from different sources, so check each against the server and client you deploy. Confirm the server version with redis-cli INFO server and read the redis_version field.
Best Value
| Item | Stated version or requirement | Source |
|---|---|---|
| Redis server for the guide’s example | Redis 7.0 or later, because the example relies on a reply shape available from that version | Redis streaming guide for redis-py |
XAUTOCLAIM |
Added in Redis 6.2 | Redis streaming guide for redis-py |
XREADGROUP |
Available since Redis Open Source 5.0.0 | XREADGROUP command reference |
| Python for the guide’s example | Python 3.9 or later | Redis streaming guide for redis-py |
redis-py for the guide’s example |
redis-py 5.0 or later | Redis streaming guide for redis-py |
XACKDEL and XDELEX |
Added in Redis 8.2, with enhanced stream operations for coordination among groups | Redis Streams overview |
| Idempotent message processing | Added in Redis 8.6, aimed at production and deduplication use | Redis Streams overview |
Do not assume the Redis 8.2 and 8.6 features exist on older servers. Code that relies on them should check the server version at startup.
WRedis: what its documentation establishes
The wredis project page on PyPI shows this interface:
from wredis.streams import RedisStreamManager
sm = RedisStreamManager(host="localhost")
sm.add_to_stream("events", {"action": "login", "user": "alice"})
@sm.on_message("events", group_name="my_group", consumer_name="worker_1")
def process(data):
print(data)
sm.wait()
The example appends an event with add_to_stream, registers a handler for the events stream under the group my_group and consumer worker_1, and then calls wait. Distinct worker processes need distinct consumer names within a group, so the hardcoded worker_1 would have to be parameterised in a real deployment. The page lists add_to_stream, on_message, exist, read_from_stream, wait, and delete_stream as its Streams methods.
The same page also documents Queue and Pub/Sub modules. Their delivery semantics differ from consumer-group semantics, so do not mix them in one pipeline on the assumption that they behave alike.
What is not established
The package page does not describe the following, and nothing in this guide verifies WRedis under failure:
Quick Recap
- whether a message is acknowledged before or after the handler runs, or whether a handler exception leaves it pending;
- whether pending entries are ever reclaimed, and with what idle threshold;
- how errors, retries, and dead letters are handled;
- how retention and trimming interact with pending entries;
- how a crashed worker’s entries are recovered, and whether that matches the flow described above.
Checks before using WRedis for reliability-critical work
- Pin the exact
wredisversion and read the Streams module source for that version, not only the project page. - Locate the acknowledgement call and confirm where it sits relative to the handler invocation.
- Send
SIGKILLto a worker in the middle of a handler, then confirm another consumer claims the entry and processes it. - Confirm whether the wrapper ever reads its own pending entries on restart, and whether it calls
XAUTOCLAIMor an equivalent. - Confirm the retention behaviour, including what happens to entries that are still pending when they are trimmed.
Which path to choose
- Choose
redis-pywith the flow in this guide when recovery is a requirement and you need explicit control over acknowledgement timing, reclaim thresholds, and retention. This is the path whose behaviour Redis documents. - Choose WRedis only after completing the checks above for the exact version you intend to deploy. Its decorator style is a reasonable fit for internal tools and prototypes, but its documentation does not yet establish the recovery behaviour a reliability-critical pipeline depends on.
Product prices and availability are accurate as of the date/time indicated and are subject to change. Any price and availability information displayed on Amazon at the time of purchase will apply.

