Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
27 commits
Select commit Hold shift + click to select a range
02c168f
Prefetch tasks in per-process batches
codingjoe Oct 1, 2026
e0654ba
Add dramatiq to the benchmark and refresh the chart
codingjoe Oct 1, 2026
dc8665f
Measure each queue at its own default read-ahead
codingjoe Oct 1, 2026
71aa089
Measure every queue at the same read-ahead
codingjoe Oct 1, 2026
bab323f
Read 128 messages ahead in the queue comparison
codingjoe Oct 1, 2026
6a13e13
Delete the dramatiq results the benchmark now writes
codingjoe Oct 1, 2026
cd36fc2
Merge remote-tracking branch 'origin/main' into codingjoe-task-prefet…
codingjoe Oct 7, 2026
f0192a8
Refresh the chart after the merge
codingjoe Oct 7, 2026
385a28a
Buffer prefetched tasks in a priority queue
codingjoe Oct 7, 2026
3c95fd7
Document the buffer's lease dwell instead of planning a fix
codingjoe Oct 7, 2026
178b7f2
Scope the wait timeout and simplify the prefetch derivation
codingjoe Oct 7, 2026
6db87c9
Drop the worker startup log
codingjoe Oct 7, 2026
e36bd01
Let the worker budget stop the fetcher on its own
codingjoe Oct 7, 2026
8d3e9f2
Order the buffered task results themselves
codingjoe Oct 7, 2026
8802097
Rewrite the new docstrings in simplified technical English
codingjoe Oct 7, 2026
981762b
Trim the prefetch section of the README
codingjoe Oct 7, 2026
804b2b2
Drop the comparison docstring
codingjoe Oct 7, 2026
8ee65d5
Move the fetch loop into its own method
codingjoe Oct 7, 2026
c8ac2f4
Hand the prefetch outcome over with one future
codingjoe Oct 7, 2026
1bd3f8a
Reduce the prefetch note to one paragraph
codingjoe Oct 7, 2026
5d46ecb
Work buffered tasks off in lease order
codingjoe Oct 7, 2026
97f4619
Trim the prefetch flag help
codingjoe Oct 7, 2026
ea9b8d2
Resolve the completion future on every exit path
codingjoe Oct 7, 2026
a96fb49
Drop the comment on the completion guard
codingjoe Oct 7, 2026
b6b6a5f
Drop the prefetch paragraph
codingjoe Oct 7, 2026
38c558d
Spawn the child in the prefetch failure tests
codingjoe Oct 7, 2026
dc33984
Refresh the comparison charts
codingjoe Oct 7, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
31 changes: 17 additions & 14 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@
<picture>
<source media="(prefers-color-scheme: dark)" srcset="https://github.com/codingjoe/threadmill/raw/main/docs/images/backend-comparison-dark.svg">
<source media="(prefers-color-scheme: light)" srcset="https://github.com/codingjoe/threadmill/raw/main/docs/images/backend-comparison-light.svg">
<img alt="Tasks per second with one worker: threadmill 5,023, celery 1,951, django-tasks-db 1,942, dramatiq 669, django-tasks-rq 80." src="https://github.com/codingjoe/threadmill/raw/main/docs/images/backend-comparison-light.svg">
<img alt="Tasks per second with one worker: threadmill 11,977, dramatiq 7,168, celery 2,183, django-tasks-db 2,154, django-tasks-rq 87." src="https://github.com/codingjoe/threadmill/raw/main/docs/images/backend-comparison-light.svg">
</picture>
</p>

Expand Down Expand Up @@ -95,12 +95,15 @@ uv run manage.py threadmill worker --max-tasks 1000 --max-tasks-jitter 100

This 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.

The limit is soft. A worker drains its buffer and the batch in hand before it stops. It can then run about twice `--prefetch-count` tasks more than the configured maximum.

Should a worker crash or be killed, the pool will automatically restart it.

#### Shutdown

A graceful shutdown is possible with `SIGTERM` or a keyboard interrupt.
All workers will finish the tasks they acquired and acknowledge them.
All workers finish the tasks they acquired and acknowledge them. This includes the tasks in their prefetch buffer.
A hard kill cannot be intercepted, so the lease reaper collects the buffered tasks after the lease expires.

You can use `--exit-empty` to exit immediately after all tasks have been processed,
which might be useful for draining a one-off queue.
Expand All @@ -125,24 +128,24 @@ uv run manage.py threadmill inspector
The `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. |
| Option | Default | Description |
| ------------------- | ------------------------- | ----------------------------------------------------------------------------------------- |
| `lease_ttl` | `timedelta(hours=1)` | Max time from acquisition to acknowledgement 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.
A claimed task whose stored payload cannot be read any more is dropped with the
read error logged. A dropped task records no result, so it leaves the inspector
and cannot be requeued.
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. The
acknowledgement of the lease holder wins: the late result of an expired attempt
is discarded.
Keep `lease_ttl` above your worst-case runtime and above the time a task waits in a
prefetch buffer. A task that outlives its lease can still run, so a retry can run
at the same time. The acknowledgement of the lease holder wins. The late result
of an expired attempt is discarded.

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
Expand Down
25 changes: 22 additions & 3 deletions benchmarks/chart.py
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,14 @@ class Theme:
accent_bar="#6366f1",
)

DIAGNOSTIC_QUEUES = frozenset({"threadmill (no prefetch)"})
"""Queues that the benchmark measures and the chart leaves out.

The harness runs threadmill twice to compare its prefetch buffer with a single
task. Against a local broker the two results differ by less than one percent. A
chart with both rows ranks them on measurement noise.
"""


@dataclasses.dataclass(frozen=True, kw_only=True, slots=True)
class QueueResult:
Expand Down Expand Up @@ -108,6 +116,8 @@ def read_results(json_path: pathlib.Path) -> list[QueueResult]:
]
results = []
for queue_name in queue_names:
if queue_name in DIAGNOSTIC_QUEUES:
continue
process = means[("test_process_queue__benchmark", queue_name)]
start = means[("test_start_worker__benchmark", queue_name)]
enqueue = means[("test_enqueue__benchmark", queue_name)]
Expand Down Expand Up @@ -143,6 +153,13 @@ def build_chart(results: list[QueueResult], theme: Theme) -> str:
"""Return the chart as an SVG document drawn in the given theme."""
fastest = results[0]
scale = PLOT_WIDTH / fastest.throughput
minimum_task_count = min(result.task_count for result in results)
maximum_task_count = max(result.task_count for result in results)
depth_label = (
f"{minimum_task_count:,}"
if minimum_task_count == maximum_task_count
else f"{minimum_task_count:,}–{maximum_task_count:,}"
)
row_centers = [
FIRST_ROW_CENTER + index * ROW_HEIGHT for index in range(len(results))
]
Expand All @@ -160,7 +177,7 @@ def build_chart(results: list[QueueResult], theme: Theme) -> str:
text(
28,
68,
f"{fastest.task_count:,} trivial tasks per queue · one worker process, "
f"{depth_label} trivial tasks per queue · one worker process, "
"one thread · higher is better",
theme=theme,
size=12.5,
Expand Down Expand Up @@ -204,8 +221,10 @@ def build_chart(results: list[QueueResult], theme: Theme) -> str:
text(
28,
footnote_y,
"threadmill, django-tasks-db and django-tasks-rq read one message at a "
"time; celery and dramatiq four.",
# joe: width checked by hand (right edge 856.1 of 900 at 11.5px); add a
# width guard if the canvas width or the font stack changes.
"One process and one thread each. Threadmill, celery and dramatiq read "
"128 ahead; django-tasks-db and -rq read one message at a time.",
theme=theme,
size=11.5,
fill=theme.faint,
Expand Down
Loading
Loading