taskiq-redis-streams
An independent, single-node Redis Streams broker for Taskiq. It uses bounded local prefetch and consumer-heartbeat pending-entry recovery.
Redis Streams delivery is at least once. A task can execute more than once after a worker crash or failed acknowledgement, so handlers must be idempotent.
Installation
uv add taskiq-redis-streams
Usage
from taskiq_redis_streams import RedisStreamsBroker
broker = RedisStreamsBroker(
"redis://localhost:6379/0",
queue_name="default",
namespace="my-service",
)
@broker.task
async def process_order(order_id: str) -> None:
# Make this operation idempotent.
...
taskiq worker my_app:broker
The consumer group always starts at 0, so tasks published before the first
worker starts are consumed.
Features
- One namespaced Redis Stream and consumer group per Taskiq queue.
- At-least-once delivery through Redis Streams consumer groups; task handlers must be idempotent.
- Bounded broker-local prefetch:
xread_countlimits each read andmax_pendingcaps fetched-but-unacknowledged entries per listener. - Consumer-heartbeat recovery: live workers can run arbitrarily long tasks without recovery being tied to a task execution timeout.
- Atomic orphan reclaim: Redis verifies that the PEL owner is unchanged and
its heartbeat is absent before
XCLAIMtransfers an entry. - A dedicated heartbeat connection, so a blocking
XREADGROUPcannot starve liveness renewal when the task delivery pool is constrained. - Fast listener-close handoff: entries still buffered inside the broker move to
an internal
abandonedconsumer and are reclaimable immediately. - Retryable Redis listener errors use exponential backoff; Taskiq cancellation is propagated so shutdown handoff still runs.
Delivery and Recovery Flow
flowchart TD
producer[Taskiq producer] -->|XADD data| stream[(Redis Stream)]
startup[Worker startup] --> group[XGROUP CREATE consumer group]
group --> listener[Broker listen loop]
listener --> heartbeat[Refresh consumer heartbeat TTL]
heartbeat --> capacity{max_pending slot available?}
capacity -->|no| wait_ack[Wait for a successful ACK]
wait_ack --> capacity
capacity -->|yes| reclaim_due{Reclaim scan due?}
reclaim_due -->|yes| pending[XPENDING RANGE]
pending --> lease{Owner heartbeat exists?}
lease -->|no| claim[Lua: verify PEL owner then XCLAIM]
lease -->|yes| read
reclaim_due -->|no| read
stream --> read[XREADGROUP new entries]
claim --> buffer[Broker-local buffer]
read --> buffer
buffer --> deliver[Yield AckableMessage to Taskiq]
deliver --> execute[Execute task]
execute --> ack[XACK when Taskiq acknowledges]
ack --> capacity
listener_close[Listener cancellation or close] --> buffered{Still in local buffer?}
buffered -->|yes| abandoned[XCLAIM to abandoned consumer]
buffered -->|already yielded| drain[Keep heartbeat until broker shutdown]
worker_loss[Crash or forced stop] --> expired[Heartbeat TTL expires]
expired --> pending
abandoned --> pending
Behavior
One broker instance serves one Taskiq queue and uses these namespaced keys:
<namespace>:stream:<queue_name>
<namespace>:workers:<queue_name>
<namespace>:heartbeat:<queue_name>:<consumer_name>
Use a distinct namespace for applications that share Redis. The default is
taskiq.
Each broker instance generates its own Redis consumer name. Consumer names are not configurable, so a restarted worker receives a new identity. Active worker consumers renew a Redis TTL heartbeat and pending entries owned by a consumer whose heartbeat has expired are eligible for recovery.
xread_count controls a single XREADGROUP batch. max_pending independently
caps entries fetched by a listener but not successfully acknowledged; both
default to 100. Once the cap is reached, the listener leaves new work for
other consumers. Set max_pending=None to disable that local cap.
max_connection_pool_size applies to task delivery commands. The broker keeps
one separate Redis connection for heartbeat renewal, so a blocking read cannot
prevent the consumer lease from being refreshed.
The broker scans the consumer group's PEL periodically. It renews each active
consumer's heartbeat every consumer_heartbeat_interval milliseconds (default
10000). A heartbeat remains live for consumer_heartbeat_ttl milliseconds
(default 30000) after its most recent successful refresh. Once that lease
expires, the next PEL scan claims entries from that consumer. XCLAIM is
guarded by an atomic Redis heartbeat check, so concurrent workers cannot both
recover the same entry. Configure a TTL longer than the renewal interval to
allow for normal scheduling and Redis latency. Reclaimed entries are handled
before new entries when a scan is due.
Taskiq's timeout label still controls task execution time, but it does not
control Redis Streams recovery. A live worker can therefore run long tasks
without their PEL entries being reclaimed solely because of task duration.
On listener close, only messages still in the broker-local buffer are handed to
an internal abandoned consumer and become reclaimable immediately. Already
yielded messages may be executing. Their worker heartbeat remains active until
broker shutdown, allowing Taskiq to drain them without a duplicate reclaim.
Retryable Redis errors while listening are retried with exponential backoff from 100 ms to 5 seconds. Cancellation from Taskiq is never retried; it propagates into the listener so the normal listener-close handoff can run.
Redis Cluster, Sentinel, delayed tasks, and a stream-level dead-letter queue are not part of the first release.
Development
Tests clear a dedicated Redis database before and after every test:
TEST_REDIS_URL=redis://127.0.0.1:7000/14 uv run pytest -q
Do not point TEST_REDIS_URL at a database containing application data.
uv sync --all-groups
uv run ruff check .
uv run mypy
uv run pytest -q
Inspiration
The Taskiq integration follows ecosystem patterns from taskiq-redis. Recovery, local prefetch, and shutdown handoff are informed by dramatiq-redis-streams.
This project is independently maintained and is not affiliated with either upstream project.
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file taskiq_redis_streams-0.1.0.tar.gz.
File metadata
- Download URL: taskiq_redis_streams-0.1.0.tar.gz
- Upload date:
- Size: 9.8 kB
- Tags: Source
- Uploaded using Trusted Publishing? Yes
- Uploaded via: twine/7.0.0 CPython/3.13.14
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
736b994e81b08ae36d38b5d5b99f1283c299a21315748ba358a4b8a339bafdbd
|
|
| MD5 |
ada4c4baa0d60085f29def5020338977
|
|
| BLAKE2b-256 |
b97993e3a23b6e2608266a12a7693cd4bfeb15b58f50d6695447656117462c60
|
Provenance
The following attestation bundles were made for taskiq_redis_streams-0.1.0.tar.gz:
Publisher:
release.yml on vvanglro/taskiq-redis-streams
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
taskiq_redis_streams-0.1.0.tar.gz -
Subject digest:
736b994e81b08ae36d38b5d5b99f1283c299a21315748ba358a4b8a339bafdbd - Sigstore transparency entry: 2333845100
- Sigstore integration time:
-
Permalink:
vvanglro/taskiq-redis-streams@2e655c809313c9768d2fd40dd1651a0a38e2c999 -
Branch / Tag:
refs/tags/v0.1.0 - Owner: https://github.com/vvanglro
-
Access:
public
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
release.yml@2e655c809313c9768d2fd40dd1651a0a38e2c999 -
Trigger Event:
release
-
Statement type:
File details
Details for the file taskiq_redis_streams-0.1.0-py3-none-any.whl.
File metadata
- Download URL: taskiq_redis_streams-0.1.0-py3-none-any.whl
- Upload date:
- Size: 11.2 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? Yes
- Uploaded via: twine/7.0.0 CPython/3.13.14
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
670edc6811e6cb487fc08170768adb5af191d3296ebe69810564aebd30616cb0
|
|
| MD5 |
db39536641234b373586075857cc63af
|
|
| BLAKE2b-256 |
5d7691ae471ac59c18bf860c5eceb73364e6ea9ae4e567f7609fe2299f752d08
|
Provenance
The following attestation bundles were made for taskiq_redis_streams-0.1.0-py3-none-any.whl:
Publisher:
release.yml on vvanglro/taskiq-redis-streams
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
taskiq_redis_streams-0.1.0-py3-none-any.whl -
Subject digest:
670edc6811e6cb487fc08170768adb5af191d3296ebe69810564aebd30616cb0 - Sigstore transparency entry: 2333845123
- Sigstore integration time:
-
Permalink:
vvanglro/taskiq-redis-streams@2e655c809313c9768d2fd40dd1651a0a38e2c999 -
Branch / Tag:
refs/tags/v0.1.0 - Owner: https://github.com/vvanglro
-
Access:
public
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
release.yml@2e655c809313c9768d2fd40dd1651a0a38e2c999 -
Trigger Event:
release
-
Statement type: