BlitzQ
A high-performance, framework-agnostic task queue for Python, built on asyncio and Redis. Part of AiNest Labs.
- Async-native workers: thousands of concurrent
async deftasks per process; sync functions on a bounded thread pool; CPU-bound functions on a process pool. - Two explicit delivery modes: reliable (Redis Streams, at-least-once, crash recovery) and fast (Redis lists, at-most-once, fewest round-trips).
- Retries with backoff and jitter, dead letters, delayed tasks, periodic tasks, timeouts, cancellation, results, multiple queues with per-queue concurrency, task priority within a queue, cross-worker rate limiting.
- Framework-agnostic core with optional FastAPI/Starlette, Django and Flask helpers.
- Measured against Celery with a reproducible benchmark suite; results below, including where BlitzQ is slower.
Status: 1.0.0, published on PyPI.
Contents
Requirements · Installation · Quick start · Tasks · Workers · Queues · Retries · Scheduling · Results · Frameworks · Configuration · Reliability · Tests · Benchmarks · Limitations · Roadmap
Requirements
- Python 3.12 or 3.13
- Redis 7.0+ (tested with 7.4). Redis Cluster is not supported.
- Runtime dependencies:
redis(redis-py ≥ 5),msgspec,typer
Installation
pip install blitzq
pip install "blitzq[django]" # optional: Django integration dependency
pip install "blitzq[flask]" # optional: Flask integration dependency
pip install "blitzq[monitoring]" # optional: Prometheus exporter
Or install from source:
git clone https://github.com/ainest-labs/BlitzQ.git && cd BlitzQ
pip install .
Start a local Redis with docker compose up -d redis. More in
docs/installation.md.
Quick start
# app.py
import asyncio
from blitzq import Queue, RetryPolicy
queue = Queue(name="default", redis_url="redis://localhost:6379/0")
@queue.task(retries=3, retry_policy=RetryPolicy(initial_delay=1, max_delay=60, backoff=2, jitter=True))
async def process_order(order_id: str) -> dict:
return {"order_id": order_id, "processed": True}
async def main():
task = await process_order.enqueue("ORD-123")
print(task.id)
result = await queue.get_result(task.id, timeout=10)
print(result) # {'order_id': 'ORD-123', 'processed': True}
if __name__ == "__main__":
asyncio.run(main())
blitzq worker app:queue # terminal 1: executes tasks
python app.py # terminal 2: enqueues and waits for the result
Runnable examples live in examples/: basic tasks, multiple queues,
retries, scheduling, FastAPI, Flask and Django.
Creating and executing tasks
@queue.task # async: runs on the worker's event loop
async def fetch(url: str) -> int: ...
@queue.task(timeout=30) # sync: runs on the worker's thread pool
def resize(path: str, width: int) -> str: ...
@queue.task(executor="process") # CPU-bound: runs on a process pool
def crunch(n: int) -> int: ...
| Call | Kind | Notes |
|---|---|---|
await task.enqueue(*args, **kwargs) |
async, non-blocking | returns a TaskHandle once Redis accepted the message |
task.enqueue_sync(*args, **kwargs) |
sync, blocking | for sync code; raises RuntimeError inside a running event loop |
await task.enqueue_many([(a,), (b,)]) / enqueue_many_sync |
async / blocking | one pipelined round-trip for many calls |
task.options(queue=, delay=, eta=, task_id=, correlation_id=, headers=, timeout=) |
- | returns a bound task with .enqueue, .enqueue_sync, .enqueue_many |
await queue.send("name", args, kwargs) / send_sync |
async / blocking | enqueue by name without importing the task code |
task(*args) |
local call | runs the function directly in the current process |
Serialization. Arguments and results are encoded with msgspec (MessagePack
by default; Serializer("json") is also available). Supported: None, bool,
int, float, str, bytes, lists, tuples, dicts and datetime. UUID,
Decimal, dataclasses and msgspec.Struct are accepted but arrive as plain data
(str/dict). Anything else is rejected at enqueue time; Serializer(enc_hook=...)
adds conversions. Data from Redis is never unpickled. Pass identifiers, not ORM
objects, requests or sessions.
Inside a task, blitzq.current_task() returns a TaskContext with id,
attempt, max_attempts, queue, correlation_id and headers. Correlation ids
and headers (e.g. traceparent) propagate automatically to tasks enqueued from
within a task.
Starting workers
Workers are separate processes. They never share memory, request context or database sessions with your web processes.
blitzq worker app:queue --queues default,emails --concurrency 200
blitzq worker app:queue --queues images --concurrency 8 --processes 8
blitzq worker app:queue --queue-concurrency images=4 --threads 32 --log-format json --metrics-port 9100
SIGTERM/SIGINT stop fetching, wait up to --shutdown-timeout seconds for
running tasks, then requeue unfinished ones; a second signal forces the stop. In
Python you can also run await Worker(queue, concurrency=50).run().
Async concurrency, threads, processes: --concurrency bounds how many tasks
run at once in one worker process. Async tasks are cheap coroutines (use
hundreds). Sync tasks occupy one of --threads pool threads. executor="process"
tasks occupy one of --processes pool processes. Scale CPU capacity by running
more worker processes. See docs/performance_tuning.md.
Multiple queues and routing
queue = Queue("default", routes={"app.images.*": "images", "*.send_*": "emails"})
@queue.task(queue="reports") # explicit queue beats routing rules
def build_report(report_id: int): ...
await build_report.options(queue="urgent").enqueue(7) # per call beats both
Each subscribed queue has its own fetch loop, and fetching reserves capacity
first, so a flooded queue cannot starve a quiet one. Cap busy queues with
--queue-concurrency, and run separate worker pools per queue to scale them
independently (tested in tests/integration/test_redis_basic.py).
Priority within a queue
@queue.task(priority="high") # decorator default
async def urgent(): ...
await task.options(priority="low").enqueue(x) # per-call override
"high"/"normal" (default)/"low". Every worker checks a queue's levels in
that order on every batch, always - there's no separate config to remember, so
a priority="high" call is never silently unheard. All three levels share the
queue's one concurrency budget (priority reorders what runs next, it doesn't
add capacity), and crash recovery/at-least-once semantics apply identically to
every level. Full detail, including the one latency trade-off it makes (a
lone low-priority message can wait up to block_timeout longer when the
queue is otherwise idle): docs/architecture.md#task-priority.
Rate limiting
@queue.task(rate_limit="10/s") # or "100/m", "1000/hour"
async def call_downstream(): ...
Caps how often the task starts, enforced with a Redis-backed token bucket
shared across every worker process - "10/s" means 10/s total, not per
worker. A task over its limit is not executed and not counted as a retry or a
failure; it's rescheduled for when a slot should be free, so a rate-limited
backlog shows up as queue depth, not as a worker sleeping.
docs/architecture.md#rate-limiting.
Retries
@queue.task(
retries=4, # up to 5 attempts
retry_policy=RetryPolicy(initial_delay=1, backoff=2, max_delay=300, jitter=True,
retry_on=(ConnectionError,), dont_retry_on=(PermissionError,)),
)
async def call_api(): ...
@queue.task(retries=10)
async def poll_export(export_id: str):
if not ready(export_id):
raise Retry(delay=30) # explicit retry; still bounded by retries=
Delays are min(max_delay, initial_delay * backoff**(n-1)), drawn from
[d/2, d] with jitter. Retries are scheduled in Redis, not slept, so workers
keep processing other tasks. Retries keep the same task id and pass through normal
queue limits. After the last attempt the task is dead-lettered. Inspect and
replay with blitzq dead-letter list and blitzq task retry ID, or
await queue.retry(id).
Scheduling
await remind.options(delay=60).enqueue("stand up") # relative
await remind.options(eta=datetime(2026, 10, 1, 9, tzinfo=UTC)).enqueue("meeting")
await handle.cancel() # certain while still scheduled
@queue.periodic("0 9 * * 1-5", tz="Europe/Berlin", missed="run_once")
async def weekday_report(): ...
@queue.periodic(Every(300)) # every 5 min, epoch-aligned
async def cleanup(): ...
Periodic tasks are dispatched by blitzq scheduler app:queue. Run two or more
for availability: each occurrence is claimed atomically, so it is dispatched once.
After downtime, missed="run_once" (default) runs the latest missed occurrence,
"run_all" runs each (up to 100) and "skip" runs none.
Results and task inspection
result = await queue.get_result(task_id, timeout=10) # raises TaskFailed / ResultTimeout
info = await queue.inspect(task_id) # TaskInfo: state, attempt, timings, error
state = await queue.status(task_id) # TaskState or None
States: scheduled, queued, running, retrying, succeeded, failed,
dead_lettered, cancelled. Final states and results are stored by default
(store_results=True, result_ttl=86400). queued/running/retrying are
recorded only with track_state=True, which costs one extra write per
transition. Disable result storage per task with @queue.task(store_result=False).
Protect result access in your application: results may be sensitive.
blitzq task inspect <id> --app app:queue
blitzq queue stats --app app:queue
Framework integrations
| Framework | Publish with | Helper |
|---|---|---|
| FastAPI, Starlette, Litestar | await task.enqueue(...) |
FastAPI(lifespan=blitzq.integrations.asgi.lifespan(queue)) |
| Django | enqueue_on_commit(task, ...) (after commit) or enqueue_sync |
blitzq.integrations.django.setup(queue) (pip install "blitzq[django]") |
| Flask | task.enqueue_sync(...) |
blitzq.integrations.flask.init_app(app, queue) (pip install "blitzq[flask]") |
| aiohttp, Sanic, Quart, scripts, notebooks | async or sync API | none needed |
Deploy web apps and workers as separate processes: web apps publish, workers execute. Details in docs/framework_integration.md.
Configuration
Queue(...) parameter |
Default | Meaning |
|---|---|---|
name |
"default" |
default queue |
redis_url |
$BLITZQ_REDIS_URL or redis://localhost:6379/0 |
rediss:// for TLS |
mode |
"reliable" |
"reliable" or "fast" |
namespace |
"blitzq" |
Redis key prefix |
store_results / result_ttl |
True / 86400 |
result storage |
track_state |
False |
record non-final states |
routes |
none | glob patterns → queue, or a callable |
default_retries / default_retry_policy / default_timeout |
0 / RetryPolicy() / none |
task defaults |
visibility_timeout |
60 |
reliable mode: seconds before an abandoned message is redelivered |
max_deliveries |
5 |
reliable mode: dead-letter after this many deliveries |
dead_letter_max |
100000 |
dead letters kept (oldest dropped) |
serializer |
Serializer() |
Serializer("json"), enc_hook=, max_message_size= |
redis_options |
none | extra redis-py connection options (TLS certs, timeouts, ...) |
Worker options: blitzq worker --help. Operational guidance, including Redis
persistence, memory and security: docs/operations.md.
Reliability modes and delivery guarantees
| reliable (default) | fast | |
|---|---|---|
| Guarantee | at-least-once | at-most-once |
| Worker killed mid-task | redelivered after visibility_timeout |
task lost |
| Graceful shutdown | unfinished tasks requeued | unfinished tasks requeued |
| Redis ops per task (no results) | fetch (batched) + ack (group-committed) | fetch (batched) |
Neither mode provides exactly-once execution. In reliable mode a task can run
twice, for example when a worker dies after finishing a task but before its ack
reached Redis, so make side effects idempotent (use current_task().id as an
idempotency key). Everything is only as durable as your Redis persistence. Full
details: docs/delivery_guarantees.md.
Running tests
pip install -e ".[dev]"
docker compose up -d redis
pytest -q # 212 tests: unit, Redis integration (both modes), crash/recovery
pytest -q tests/unit # no Redis needed
ruff check src tests && mypy
Benchmarks
The suite in benchmarks/ runs identical workloads against
BlitzQ, Celery 5.6 and Huey 3.4 on the same Redis and reports raw per-task data,
medians over 5 repetitions, variability and resource usage. Methodology, profiles
and the full results are in docs/benchmarking.md and
benchmarks/results/published/.
Results (Linux container, 12 CPUs, Redis 7.4, Celery 5.6.3, 5 repetitions, medians)
Profile A = equal settings: the same sync task code, 16 execution slots, one worker process, one producer making one sync call per task. Profile B = each system tuned: 4 worker and 4 producer processes each. All rows compare matched delivery guarantees (at-least-once unless marked early-ack).
| Workload | Profile | BlitzQ tasks/s | Celery tasks/s | Ratio |
|---|---|---|---|---|
| no-op, producer and workers concurrent | A | 2,604 | 855 | 3.0x |
| no-op, pre-loaded backlog (worker capacity) | A | 5,503 | 967 | 5.7x |
| no-op, pre-loaded backlog | B | 70,696 | 3,738 | 18.9x |
| 1 KiB payload | B | 59,270 | 3,129 | 18.9x |
| 100 KiB payload | B | 3,812 | 1,265 | 3.0x |
| 20 ms I/O task | A | 755 | 720 | 1.05x |
| 20 ms I/O task (BlitzQ async, Celery prefork) | B | 41,807 | 2,632 | 15.9x |
| retry-heavy (50% fail once) | B | 8,577 | 1,370 | 6.3x |
| result storage on | B | 49,521 | 2,748 | 18.0x |
| CPU-bound (threads vs prefork) | A | 73 | 486 | 0.15x |
| CPU-bound (process pools, 8 each) | B | 472 | 445 | 1.06x |
Other measured results:
- Where BlitzQ is worse. CPU-bound work on BlitzQ's default thread executor is
GIL-bound (0.15x); use
executor="process". Recovery after all Redis connections were dropped took 5.1 s vs 4.4 s for Celery. With tuned settings, delayed tasks started 13 ms late at p50 vs 0.6 ms for Celery, which keeps ETA tasks in worker memory. The scheduled workload also costs BlitzQ more Redis commands per task (11–16 vs 10) because of promoter polling. - Worker crash (SIGKILL at 30%, at-least-once). 0 tasks lost in both systems. The last task completed 16.9 s after the kill with BlitzQ vs 101.6 s with Celery, whose Redis transport restores unacknowledged messages on a ~100 s cycle. BlitzQ re-executed 1 task across 5 runs.
- Early-ack crash. Both lost 15–16 of 2,000 tasks per crash, the tasks that were executing.
- Redis connections dropped (at-least-once). 0 lost in both. Duplicate executions: 0 for BlitzQ, 77 over 5 runs for Celery.
- Footprint. Peak worker RSS was 39 MB for BlitzQ (1 process, 16 threads) vs about 780 MB for Celery (prefork, 16 children). Worker CPU per task was 2.7–21x lower for BlitzQ on the no-op, payload, burst, retry and result-storage workloads.
These are single-machine numbers with trivial task bodies; real tasks that do meaningful work shrink the differences. See the full report, which lists every run, variability, settings and limitations.
BlitzQ vs Huey (same machine, 5 repetitions, medians)
Huey has no at-least-once mode of its own - it pops a message before executing it, like Celery's early-ack mode or BlitzQ's fast mode - so it is only compared against those, never against BlitzQ's reliable mode.
| Workload | Profile | BlitzQ tasks/s | Huey tasks/s | Ratio |
|---|---|---|---|---|
| no-op | A (equal, 16 threads each) | 2,847 | 1,499 | 1.9x |
| no-op | B (tuned, 4 processes each) | 126,462 | 6,024 | 21.0x |
| 20 ms I/O task | A | 767 | 762 | 1.01x |
| 20 ms I/O (BlitzQ async, Huey greenlet) | B | 62,120 | 8,657 | 7.2x |
| retry-heavy | A | 1,914 | 609 | 3.1x |
| CPU-bound (threads vs threads) | A | 73 | 70 | 1.0x |
| CPU-bound (process pools, 8 each) | B | 467 | 478 | 0.98x |
- Worker crash (early-ack). Both lost 15–16 of 2,000 tasks per crash (the tasks executing at the kill), and both recovered in about 5 s.
- Footprint. Huey used 4 Redis commands per task on every workload (vs BlitzQ's 1–7, workload-dependent) but had similar or higher worker CPU per task than BlitzQ on most workloads.
- Scheduled tasks started 1.0–2.1 s late at p50 on a 2 s delay (Huey's default 1 s scheduler poll interval, left at its shipped default).
- A bug found while building this comparison: Huey's
greenletworker type does not monkey-patch the standard library itself, so a synchronous task usingtime.sleepblocked every greenlet on that worker instead of yielding, measuring 48 tasks/s on a 20 ms I/O workload untilgevent.monkey.patch_all()was applied before Redis/Huey are imported (benchmarks/apps/huey_app.py); afterwards it measured 2,750 tasks/s on the same workload.
Full report: benchmarks/results/published/huey_suite/REPORT.md.
Reproduce:
docker compose up -d redis redis-stats
docker compose --profile bench run --rm bench python -m benchmarks.runner --config benchmarks/configs/suite.yaml
docker compose --profile bench run --rm bench python -m benchmarks.runner --config benchmarks/configs/huey_suite.yaml
python -m benchmarks.report --input benchmarks/results/<run-dir>
Current limitations
- Redis only (plus an in-memory test broker). Redis Cluster is not supported. Valkey, Dragonfly and KeyDB are untested.
- No exactly-once execution. Reliable mode is at-least-once; fast mode loses in-flight tasks when a worker process dies.
- Thread and process tasks cannot be interrupted on timeout or cancellation.
- Cancelling a queued or running task is best effort (revocation synced about every second). Only scheduled tasks are cancelled with certainty.
- Without
track_state, queued and running tasks have no inspectable record. - Task chains/groups/chords and a web dashboard are not implemented. Priority and rate limiting within a queue are (see above).
- Delayed-task precision is bounded by the promoter poll interval (0.5 s default) for tasks scheduled earlier than anything already pending.
See docs/roadmap.md for why Redis Cluster support isn't
there yet, and what's already been fixed (get_result() used to poll; it's
now notification-driven).
License
MIT, see LICENSE.
Built by AiNest Labs. Issues and pull requests welcome - see CONTRIBUTING.md.
Metadata
Release files for blitzq 1.0.1
For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.
Source distribution (sdist)
| File | Size | Uploaded | |
|---|---|---|---|
| blitzq-1.0.1.tar.gz | 91.7 kB | Details |
Built distribution (wheel)
| File | Interpreter | ABI | Platform | Reset |
|---|---|---|---|---|
| blitzq-1.0.1-py3-none-any.whl | Python 3 | none | any | Details |
Total release size: 170.6 kB
Release files / blitzq-1.0.1.tar.gz
| Download URL | blitzq-1.0.1.tar.gz |
|---|---|
| Size | 91.7 kB |
| Tags | Source |
|
SHA-256 checksum How to use checksums |
cd9fb4cccdb510f53eb5e9ab4c0902d524355a8a0232a69ad8300b95768cfc13
|
|
BLAKE2b-256 checksum How to use checksums |
9d758bda91217e3e8209c5b05c82fe2a748e88c132a74ddb230f51b75671d3ab
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
Yes |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|
Provenance
Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.
PyPI Publish Attestation
PyPI verified that this artifact, at this checksum, originated from the publisher listed below.
Signed by GitHub Actions, verified by PyPI on Sep 30, 2026.
Transparency logRelease files / blitzq-1.0.1-py3-none-any.whl
| Download URL | blitzq-1.0.1-py3-none-any.whl |
|---|---|
| Size | 78.9 kB |
| Tags | Python 3 |
|
SHA-256 checksum How to use checksums |
55c60e70dbba163327ba503c9f052c9c2590c737c2608a118f7960d08f36a767
|
|
BLAKE2b-256 checksum How to use checksums |
53dd6db70f2efcb6da269295487e88fe7ac7b6c403785324e96de2f897b3909f
|
| Upload date | |
|
Uploaded using Trusted Publishing? What is trusted publishing? |
Yes |
| Uploaded via |
twine/7.0.0 CPython/3.13.14
|
Provenance
Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.
PyPI Publish Attestation
PyPI verified that this artifact, at this checksum, originated from the publisher listed below.
Signed by GitHub Actions, verified by PyPI on Sep 30, 2026.
Transparency log