Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
25 changes: 25 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,31 @@ jobs:
- uses: codecov/codecov-action@v7
with:
token: ${{ secrets.CODECOV_TOKEN }}
pytest-free-threaded:
name: Pytest free-threaded
strategy:
matrix:
os:
- "ubuntu-latest"
python-version:
- "3.14t"
django-version:
- "6.1"
runs-on: ${{ matrix.os }}
services:
redis:
image: redis
ports:
- 6379:6379
options: --entrypoint redis-server
env:
REDIS_URL: redis:///0
steps:
- uses: actions/checkout@v7
- uses: astral-sh/setup-uv@v7
with:
python-version: ${{ matrix.python-version }}
- run: uv run --with django~=${{ matrix.django-version }}.0 pytest -m "not benchmark"
pytest-extras:
name: Pytest
strategy:
Expand Down
2 changes: 2 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -13,9 +13,11 @@ All commands run via `uv`:
```bash
uv run pytest # full suite (incl. benchmarks, coverage)
uv run pytest -m "not benchmark" # what CI runs by default
uv run --python 3.14t pytest -m "not benchmark" # free-threaded build (uv installs 3.14t on demand)
uv run pytest -m integration # integration tests only
uv run pytest -m "integration and benchmark"
uv run pytest --benchmark-compare # compare vs main baseline (run main first)
uv run pytest benchmarks/test_scaling.py -m benchmark # process vs thread scaling; run on 3.14 and 3.14t
uvx prek run --all-files
uv run manage.py threadmill worker # run the worker pool
uv run manage.py threadmill inspector # launch the textual TUI inspector
Expand Down
19 changes: 15 additions & 4 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 9,960, dramatiq 7,177, huey 6,038, celery 2,280, django-tasks-db 2,065, django-tasks-rq 82." src="https://github.com/codingjoe/threadmill/raw/main/docs/images/backend-comparison-light.svg">
<img alt="Tasks per second with one worker process each: Threadmill 19,954, Threadmill (GIL) 9,215, dramatiq 7,251, huey 6,254, celery 2,497, django-tasks-db 1,814, django-tasks-rq 84." src="https://github.com/codingjoe/threadmill/raw/main/docs/images/backend-comparison-light.svg">
</picture>
</p>

Expand Down Expand Up @@ -78,13 +78,24 @@ The workers are inspired by Gunicorn, and the CLI is very similar.

#### Utilization

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.
Depending on your workload, you can tweak the number of worker processes and the number of threads per process. Processes always run in parallel, because each one has its own interpreter. Threads share memory, so they are cheap, but whether they run in parallel depends on your interpreter:

- On a regular GIL build, Python runs in one thread at a time. Threads still overlap waiting for IO, but CPU-bound tasks do not speed up. Raising the thread count far above the queue depth can also leave processes idle, because each process prefetches a batch sized by its thread count.
- On a [free-threaded build](https://peps.python.org/pep-0779/), threads run truly in parallel, so CPU-bound tasks scale across cores without the cost of extra processes.

A pool on a free-threaded interpreter therefore reaches the same throughput with one process running many threads as with many processes running one thread each, while using a single Redis connection pool and one copy of your application state:

```console
uv run manage.py threadmill worker --workers 4 --threads 2
uv run manage.py threadmill worker --workers 1 --threads 8
```

Threads all live in one process, so a crash or a `--max-tasks` recycle takes every thread down at once. One process per core, the default, confines that to a single process. Prefer threads when memory matters more than that isolation, and keep the default when it does not.

A free-threaded interpreter can be slower than a regular one for single-threaded work, so a GIL build stays the better choice for IO-bound tasks. Pick the interpreter that fits your workload rather than assuming free threading is an upgrade.

> [!WARNING]
> A C extension that does not declare free-threading support re-enables the GIL for the whole process, which silently costs you all thread parallelism. `hiredis` is a common offender: installing `redis[hiredis]` enables it on import. The worker logs a warning when it detects that the GIL was re-enabled.

#### Health

If your tasks leak memory, you can recycle (restart) the workers after a certain number of tasks have been processed:
Expand Down
9 changes: 9 additions & 0 deletions benchmarks/celery_app.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,8 @@
import redis
from celery import Celery

from benchmarks import cpu_work

REDIS_URL = os.environ.get("REDIS_URL", "redis://localhost:6379/0")

PROCESSED_KEY = "benchmark:processed"
Expand All @@ -31,6 +33,13 @@ def celery_echo(value):
return value


@celery_app.task(queue=cpu_work.CPU_QUEUE)
def celery_compute(value):
"""Consume CPU, then record that this task finished."""
cpu_work.count_primes()
client.incr(cpu_work.COMPLETION_KEY)


@celery_app.task
def celery_mark_processed():
"""Record that every earlier task in the queue has been processed."""
Expand Down
105 changes: 72 additions & 33 deletions benchmarks/chart.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,8 +2,13 @@

Generate the numbers first, then the chart:

uv run pytest benchmarks -m benchmark --benchmark-json=benchmark.json
uv run python benchmarks/chart.py benchmark.json
uv run pytest benchmarks/test_backends.py -m benchmark --benchmark-json=benchmark.json
uv run --python 3.14t pytest benchmarks/test_backends.py -m benchmark -k free \
--benchmark-json=benchmark-free-threaded.json
uv run python benchmarks/chart.py benchmark.json benchmark-free-threaded.json

The first run measures every queue on the GIL build. The second one supplies the
free-threading row, which the benchmark lists only where threads run in parallel.
"""

import dataclasses
Expand All @@ -23,8 +28,7 @@

FIRST_ROW_CENTER = 110
ROW_HEIGHT = 40
FOOTNOTE_GAP = 38
BOTTOM_PADDING = 64
BOTTOM_PADDING = 52

FONT = (
'system-ui, -apple-system, "Segoe UI", Roboto, "Helvetica Neue", Arial, sans-serif'
Expand All @@ -47,9 +51,6 @@ class Theme:
muted: str
"""Secondary text, and the bars that are not highlighted."""

faint: str
"""Footnotes."""

accent: str
"""The highlighted queue and its value."""

Expand All @@ -62,7 +63,6 @@ class Theme:
border="#e5e7eb",
ink="#111827",
muted="#6b7280",
faint="#9ca3af",
accent="#4f46e5",
accent_bar="#4f46e5",
)
Expand All @@ -72,19 +72,33 @@ class Theme:
border="#30363d",
ink="#e6edf3",
muted="#8b949e",
faint="#6e7681",
accent="#818cf8",
accent_bar="#6366f1",
)

DIAGNOSTIC_QUEUES = frozenset({"threadmill (no prefetch)"})
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.
"""

FREE_THREADING_QUEUE = "Threadmill (free threading)"
"""Queue the benchmark lists only where threads really run in parallel."""

DISPLAY_NAMES = {
"Threadmill (free threading)": "Threadmill",
"Threadmill": "Threadmill (GIL)",
}
"""Chart labels for the threadmill rows, keyed by benchmark queue name.

The benchmark names identify a row across runs, so they stay as they are and the
chart translates them for the reader. The free-threading result is threadmill's
headline row and carries the plain name, and the row measured next to the other
queues is labelled as the GIL baseline.
"""


@dataclasses.dataclass(frozen=True, kw_only=True, slots=True)
class QueueResult:
Expand Down Expand Up @@ -130,9 +144,37 @@ def read_results(json_path: pathlib.Path) -> list[QueueResult]:
task_count=process["extra_info"]["tasks"],
)
)
return rank(results)


def rank(results: list[QueueResult]) -> list[QueueResult]:
"""Return the results, fastest first."""
return sorted(results, key=lambda result: result.throughput, reverse=True)


def with_free_threading_result(
results: list[QueueResult], free_threading_results: list[QueueResult]
) -> list[QueueResult]:
"""Return the results with the free-threading row taken from another run.

The free-threading row is measured on a free-threaded interpreter, while every
other row is measured on the GIL build the chart ranks, so the two runs are
merged. The remaining rows of the second run repeat queues that the first run
already measured and are ignored.
"""
free_threading = [
result
for result in free_threading_results
if result.name == FREE_THREADING_QUEUE
]
if not free_threading:
raise SystemExit(f"no {FREE_THREADING_QUEUE} row in the free-threaded run")
return rank(
[result for result in results if result.name != FREE_THREADING_QUEUE]
+ free_threading
)


def text(x, y, content, *, theme, size=13, fill=None, weight=400, anchor="start"):
"""Render one SVG text element."""
return (
Expand All @@ -142,10 +184,15 @@ def text(x, y, content, *, theme, size=13, fill=None, weight=400, anchor="start"
)


def display_name(queue_name: str) -> str:
"""Return the chart label of a queue."""
return DISPLAY_NAMES.get(queue_name, queue_name)


def describe(results: list[QueueResult]) -> str:
"""Return a sentence describing the throughput of every queue."""
return "Tasks per second with one worker: " + ", ".join(
f"{result.name} {result.throughput:,.0f}" for result in results
return "Tasks per second with one worker process each: " + ", ".join(
f"{display_name(result.name)} {result.throughput:,.0f}" for result in results
)


Expand All @@ -163,8 +210,7 @@ def build_chart(results: list[QueueResult], theme: Theme) -> str:
row_centers = [
FIRST_ROW_CENTER + index * ROW_HEIGHT for index in range(len(results))
]
footnote_y = row_centers[-1] + FOOTNOTE_GAP
height = footnote_y + BOTTOM_PADDING
height = row_centers[-1] + BOTTOM_PADDING

parts = [
f'<svg xmlns="http://www.w3.org/2000/svg" viewBox="0 0 {WIDTH} {height}" '
Expand All @@ -177,8 +223,8 @@ def build_chart(results: list[QueueResult], theme: Theme) -> str:
text(
28,
68,
f"{depth_label} trivial tasks per queue · one worker process, "
"one thread · higher is better",
f"{depth_label} trivial tasks per queue · one worker process each · "
"higher is better",
theme=theme,
size=12.5,
fill=theme.muted,
Expand All @@ -192,7 +238,7 @@ def build_chart(results: list[QueueResult], theme: Theme) -> str:
text(
LABEL_X,
center + 5,
result.name,
display_name(result.name),
theme=theme,
size=14,
weight=700 if is_fastest else 400,
Expand All @@ -217,32 +263,25 @@ def build_chart(results: list[QueueResult], theme: Theme) -> str:
)
)

parts.append(
text(
28,
footnote_y,
# joe: width checked by hand (right edge 863 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, -rq and huey read one task at a time.",
theme=theme,
size=11.5,
fill=theme.faint,
)
)
parts.append("</svg>")
return "\n".join(parts)


if __name__ == "__main__":
if len(sys.argv) != 2:
raise SystemExit(f"usage: {sys.argv[0]} <benchmark.json>")
if len(sys.argv) not in {2, 3}:
raise SystemExit(
f"usage: {sys.argv[0]} <benchmark.json> [<free-threaded-benchmark.json>]"
)
queues = read_results(pathlib.Path(sys.argv[1]))
if len(sys.argv) == 3:
queues = with_free_threading_result(
queues, read_results(pathlib.Path(sys.argv[2]))
)
for theme, path in ((LIGHT_THEME, LIGHT_THEME_PATH), (DARK_THEME, DARK_THEME_PATH)):
path.write_text(build_chart(queues, theme) + "\n")
print(f"wrote {path} ({path.stat().st_size} bytes)")
for queue in queues:
print(
f"{queue.name:20s} {queue.throughput:8,.0f}/s "
f"{display_name(queue.name):20s} {queue.throughput:8,.0f}/s "
f"(enqueue {1 / queue.enqueue_seconds:8,.0f}/s, start {queue.start_seconds:.4f}s)"
)
52 changes: 52 additions & 0 deletions benchmarks/cpu_work.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
"""The shared contract every framework in the parallelism benchmark follows.

The Celery and dramatiq worker CLIs import their applications without setting up
Django, so this module must not import Django or any Django application. The work
matches ``tests.testapp.tasks.compute_workload``, which is the equivalent Django
task the threadmill arm of the benchmark runs.
"""

PRIME_TARGET = 100_000
"""How many primes to count, about one second of work on a modern core."""

COMPLETION_KEY = "benchmark:cpu_done"
"""Key each CPU task increments when it finishes.

A sentinel task cannot prove a drain on a pool with more than one thread: a free
thread can take the sentinel before its predecessors finish. Counting completions
proves it whatever the completion order.
"""

CPU_QUEUE = "cpu_benchmark"
"""Queue carrying only the CPU workload.

Defined once here because both framework applications must publish to it and the
benchmark must listen on it. A queue of its own keeps the measurement honest: a
worker process left behind by another test or worktree sharing the same Redis
instance listens on the framework's default queue and would otherwise steal tasks
and report a drain that never did the work.
"""


def count_primes(target: int = PRIME_TARGET) -> int:
"""Count the first ``target`` primes and return the count."""

def is_prime(number: int) -> bool:
if number < 2:
return False
if number in (2, 3):
return True
if number % 2 == 0:
return False
for divisor in range(3, int(number**0.5) + 1, 2):
if number % divisor == 0:
return False
return True

prime_count = 0
number = 2
while prime_count < target:
if is_prime(number):
prime_count += 1
number += 1
return prime_count
9 changes: 9 additions & 0 deletions benchmarks/dramatiq_app.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,8 @@
from dramatiq.results import Results
from dramatiq.results.backends.redis import RedisBackend

from benchmarks import cpu_work

REDIS_URL = os.environ.get("REDIS_URL", "redis://localhost:6379/0")

PROCESSED_KEY = "benchmark:processed"
Expand Down Expand Up @@ -47,3 +49,10 @@ def dramatiq_echo(value):
def dramatiq_mark_processed():
"""Record that every earlier task in the queue has been processed."""
client.incr(PROCESSED_KEY)


@dramatiq.actor(broker=redis_broker, queue_name=cpu_work.CPU_QUEUE)
def dramatiq_compute(value):
"""Consume CPU, then record that this task finished."""
cpu_work.count_primes()
client.incr(cpu_work.COMPLETION_KEY)
Loading
Loading