Documentation |
Issues |
Changelog |
Funding π
Durable high-performance backend for Django's task framework.
- Durability β Recover from any failures, even poorly written tasks.
- Consistency β Never lose data, even if someone unplugs the power or network.
- Utilization β Keep the CPU saturated with tasks, not with idle time or waiting for locks.
You need to have Django's Task framework set up properly.
uv add "threadmill[redis]"Add threadmill to your INSTALLED_APPS in settings.py
and configure the task backend:
# settings.py
import os
INSTALLED_APPS = [
"threadmill",
# ...
]
TASKS = {
"default": {
"BACKEND": "threadmill.backends.redis.RedisTaskBackend",
"REDIS_URL": os.getenv("REDIS_URL", "redis://localhost:6379/0"),
},
# ...
}Optionally, install the inspector dependency if you want the TUI:
uv add "threadmill[inspector]"Then launch the worker pool:
uv run manage.py threadmill workerThe workers are inspired by Gunicorn, and the CLI is very similar.
Depending on your workload, you can tweak the number of processes and threads. Processes allow for parallel compute (no GIL) while threads are great for low-memory concurrent IO.
uv run manage.py threadmill worker --workers 4 --threads 2If your tasks leak memory, you can recycle (restart) the workers after a certain number of tasks have been processed:
uv run manage.py threadmill worker --max-tasks 1000 --max-tasks-jitter 100This will restart the workers after 1000 tasks have been processed, with a random jitter of up to 100 tasks to avoid all workers restarting at the same time.
Should a worker crash or be killed, the pool will automatically restart it.
A graceful shutdown is possible with SIGTERM or a keyboard interrupt.
All workers will finish the tasks they acquired and acknowledge them.
You can use --exit-empty to exit immediately after all tasks have been processed,
which might be useful for draining a one-off queue.
The optional TUI inspector lets you watch queues, tasks, and task details in real-time.
Install it with the inspector extra and launch it from a separate terminal:
uv add "threadmill[inspector]"
uv run manage.py threadmill inspectorThe RedisTaskBackend accepts the following options under OPTIONS in your
TASKS configuration:
| Option | Default | Description |
|---|---|---|
lease_ttl |
timedelta(hours=1) |
Max processing time before the task is retried or marked FAILED. |
result_ttl |
timedelta(days=1) |
How long task results are retained before automatic removal. |
broker_interval |
timedelta(seconds=1) |
Interval between background broker maintenance passes. |
batch_size |
100 |
Max tasks to move or reap per broker pass. |
poll_interval |
timedelta(seconds=0.01) |
Base wait between idle acquire attempts, doubled after each empty poll. |
poll_max_interval |
timedelta(seconds=1) |
Max wait between idle acquire attempts. |
A task whose lease expired reaches the retry callback as an
AcknowledgementTimeout error, or is marked FAILED when nothing retries it.
Keep lease_ttl above your worst-case runtime: a task that outlives its lease
can still be running, so a retry may execute concurrently with it.
All keys for one backend alias share a Redis Cluster hash tag ({alias}), so
every multi-key operation β including the cross-queue acquire β runs on a single
shard. Scale horizontally by running additional backend aliases, not by relying
on cross-slot operations.
Pass a retry callback to @task() to retry failed tasks with a delay.
The callback receives a TaskContext β use context.attempt for the current
attempt count and context.task_result.errors[-1] for the latest error.
Return a timedelta to schedule the next attempt, or None to stop retrying.
Failed tasks are re-queued preserving their ID and error history; the broker promotes them back to the ready queue once the delay elapses.
threadmill.retry.ExponentialBackoff provides a serializable exponential
backoff strategy out of the box. It caps the delay at max_delay, stops
after max_retries attempts, and only retries exceptions listed in
expected_exceptions.
import datetime
from django.tasks import task
from requests import HTTPError
from threadmill.exceptions import AcknowledgementTimeout
from threadmill.retry import ExponentialBackoff
@task(
retry=ExponentialBackoff(
base_delay=datetime.timedelta(seconds=1),
max_delay=datetime.timedelta(minutes=5),
factor=2.0,
max_retries=5,
expected_exceptions=(HTTPError, AcknowledgementTimeout),
)
)
def fetch_github_api(url: str): ...For cases that need logic beyond what ExponentialBackoff supports,
write a callable that accepts a TaskContext and returns a timedelta
or None. Use TaskError.exception_class to filter by exception type:
import datetime
from django.tasks import task
from django.tasks.base import TaskContext
from requests import HTTPError
def retry_on_rate_limit(context: TaskContext) -> datetime.timedelta | None:
"""Retry HTTP 429 responses with exponential backoff, up to 5 attempts."""
if context.attempt >= 5:
return None
error = context.task_result.errors[-1]
if not issubclass(error.exception_class, HTTPError):
return None
return min(
datetime.timedelta(seconds=2**context.attempt),
datetime.timedelta(seconds=60),
)
@task(retry=retry_on_rate_limit)
def fetch_github_api(url: str): ...