From e2c19e194cfb0d7a3f80b95f1fe653c76dd98859 Mon Sep 17 00:00:00 2001
From: Johannes Maron
Date: Thu, 1 Oct 2026 16:16:04 +0200
Subject: [PATCH 1/7] Support free-threaded Python builds
On Python 3.14's free-threaded build one process running many threads
reaches the throughput of many processes running one thread each, while
using a single Redis connection pool and one copy of the application
state.
- Give each worker thread its own idle backoff and rotation offset. The
previous instance-wide state was read and written by every thread in
the process, so threads reset and doubled each other's delay on top of
racing on the counter.
- Warn at startup when a free-threaded interpreter runs with the GIL
enabled. A C extension that has not declared free-threading support,
hiredis being the common one, re-enables the GIL process-wide without
any other signal, so the parallelism is silently lost.
- Add a CPU-bound scaling benchmark comparing process and thread
parallelism. It shows threads give exactly 1.00x on a GIL build and
2.83x on a free-threaded build.
- Run the test suite on 3.14t in CI and declare the free-threading
classifiers.
---
.github/workflows/ci.yml | 26 +++++++
AGENTS.md | 2 +
README.md | 15 +++-
benchmarks/test_scaling.py | 129 +++++++++++++++++++++++++++++++++++
pyproject.toml | 2 +
tests/backends/test_redis.py | 60 +++++++++++++---
tests/test_executor.py | 94 +++++++++++++++++++++++++
tests/testapp/tasks.py | 7 ++
threadmill/backends/redis.py | 41 +++++++++--
threadmill/executor.py | 29 ++++++++
10 files changed, 386 insertions(+), 19 deletions(-)
create mode 100644 benchmarks/test_scaling.py
diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml
index 58c8d91..77fec46 100644
--- a/.github/workflows/ci.yml
+++ b/.github/workflows/ci.yml
@@ -43,6 +43,32 @@ 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 }}
+ # The default job reports coverage; this one proves free-threaded correctness.
+ - run: uv run --with django~=${{ matrix.django-version }}.0 pytest -m "not benchmark"
pytest-extras:
name: Pytest
strategy:
diff --git a/AGENTS.md b/AGENTS.md
index dc90695..8ed2f8a 100644
--- a/AGENTS.md
+++ b/AGENTS.md
@@ -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
diff --git a/README.md b/README.md
index 13b8506..6f99bc0 100644
--- a/README.md
+++ b/README.md
@@ -78,13 +78,22 @@ 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.
+- 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 --processes 4 --threads 2
+uv run manage.py threadmill worker --workers 1 --threads 8
```
+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:
diff --git a/benchmarks/test_scaling.py b/benchmarks/test_scaling.py
new file mode 100644
index 0000000..9ae207d
--- /dev/null
+++ b/benchmarks/test_scaling.py
@@ -0,0 +1,129 @@
+"""Benchmarks comparing process and thread parallelism for CPU-bound tasks.
+
+The backend comparison measures queue overhead with a trivial echo task, which
+says nothing about the parallelism a worker pool actually delivers. Only a
+CPU-bound task can show that, so this module drains the same fixed workload with
+several worker process and thread counts and reports how long each takes.
+
+Run both interpreters and compare the ``(1, 1)`` baseline against the threaded
+configurations:
+
+ uv run pytest benchmarks/test_scaling.py -m benchmark --benchmark-json=scaling.json
+ uv run --python 3.14t pytest benchmarks/test_scaling.py -m benchmark --benchmark-json=scaling-free-threaded.json
+
+On a GIL build extra threads cannot shorten the drain, because only one thread
+runs Python at a time. On a free-threaded build they can. Every measurement
+includes the worker pool's fixed start cost of one to two seconds, which is a
+larger share of the faster configurations, so compare each configuration against
+the ``(1, 1)`` baseline instead of reading the times as pure throughput.
+"""
+
+import dataclasses
+import sys
+
+import pytest
+from django.core.management import call_command
+from django.tasks import (
+ DEFAULT_TASK_BACKEND_ALIAS,
+ task_backends,
+)
+
+from tests.testapp.tasks import compute_workload
+from threadmill.executor import is_free_threaded_build, is_gil_enabled
+
+TASK_COUNT = 16
+"""CPU-bound tasks drained per measurement, about one second of work each."""
+
+MEASUREMENT_ROUNDS = 2
+"""Repeats per configuration, reduced to the median by the benchmark plugin."""
+
+
+@dataclasses.dataclass(frozen=True, kw_only=True, slots=True)
+class WorkerParallelism:
+ """A worker process count and a worker thread count to measure."""
+
+ workers: int
+ threads: int
+
+ @property
+ def label(self) -> str:
+ """Return a short identifier for benchmark reports."""
+ return f"{self.workers}p-{self.threads}t"
+
+
+PARALLELISMS_UNDER_TEST = (
+ WorkerParallelism(workers=1, threads=1),
+ WorkerParallelism(workers=4, threads=1),
+ WorkerParallelism(workers=1, threads=4),
+ WorkerParallelism(workers=4, threads=4),
+)
+"""Configurations spanning process-only, thread-only and mixed parallelism."""
+
+
+@dataclasses.dataclass(frozen=True, kw_only=True, slots=True)
+class CpuWorkload:
+ """A fixed CPU-bound workload processed by one worker pool configuration."""
+
+ parallelism: WorkerParallelism
+ task_count: int = TASK_COUNT
+
+ def drain(self) -> None:
+ """Queue the workload, then process it with the configured worker pool."""
+ queued_task = compute_workload.using(backend=DEFAULT_TASK_BACKEND_ALIAS)
+ for _ in range(self.task_count):
+ queued_task.enqueue()
+ call_command(
+ "threadmill",
+ "worker",
+ backend=DEFAULT_TASK_BACKEND_ALIAS,
+ queues=[compute_workload.queue_name],
+ workers=self.parallelism.workers,
+ threads=self.parallelism.threads,
+ exit_empty=True,
+ verbosity=0,
+ )
+
+
+@pytest.fixture
+def empty_queues():
+ """Delete queued tasks before and after each measurement."""
+ client = task_backends[DEFAULT_TASK_BACKEND_ALIAS].client
+
+ def delete_queued_tasks() -> None:
+ if keys := client.keys("threadmill:*"):
+ client.delete(*keys)
+
+ delete_queued_tasks()
+ yield
+ delete_queued_tasks()
+
+
+class TestThreadScaling:
+ """Measure how long a fixed CPU-bound workload takes to drain."""
+
+ @pytest.mark.benchmark
+ @pytest.mark.django_db(transaction=True)
+ @pytest.mark.parametrize(
+ "parallelism",
+ PARALLELISMS_UNDER_TEST,
+ ids=lambda parallelism: parallelism.label,
+ )
+ def test_drain_cpu_workload__benchmark(self, benchmark, parallelism, empty_queues):
+ """Benchmark the time to drain the workload with one pool configuration."""
+ benchmark.extra_info.update(
+ {
+ "workers": parallelism.workers,
+ "threads": parallelism.threads,
+ "tasks": TASK_COUNT,
+ "python": sys.version.split()[0],
+ "free_threaded_build": is_free_threaded_build(),
+ "gil_enabled": is_gil_enabled(),
+ }
+ )
+
+ benchmark.pedantic(
+ CpuWorkload(parallelism=parallelism).drain,
+ rounds=MEASUREMENT_ROUNDS,
+ iterations=1,
+ warmup_rounds=0,
+ )
diff --git a/pyproject.toml b/pyproject.toml
index 340e138..f0e7eb9 100644
--- a/pyproject.toml
+++ b/pyproject.toml
@@ -25,6 +25,8 @@ classifiers = [
"Programming Language :: Python :: 3 :: Only",
"Programming Language :: Python :: 3.14",
"Programming Language :: Python :: 3.15",
+ "Programming Language :: Python :: Free Threading",
+ "Programming Language :: Python :: Free Threading :: 3 - Stable",
"Topic :: Communications :: Email",
"Topic :: Software Development",
"Topic :: Text Processing :: Markup :: Markdown",
diff --git a/tests/backends/test_redis.py b/tests/backends/test_redis.py
index 4cb58b9..6e5e1de 100644
--- a/tests/backends/test_redis.py
+++ b/tests/backends/test_redis.py
@@ -3,6 +3,7 @@
import datetime
import logging
import queue
+import threading
import time
import typing
from dataclasses import replace
@@ -20,7 +21,11 @@
QueueRates,
QueueStats,
)
-from threadmill.backends.redis import RedisBroker, RedisTaskBackend # noqa: E402
+from threadmill.backends.redis import ( # noqa: E402
+ IdleBackoff,
+ RedisBroker,
+ RedisTaskBackend,
+)
TELEMETRY_INTERVAL = datetime.timedelta(seconds=60)
@@ -1156,7 +1161,7 @@ def test_acquire__rebase_rotation_index_beyond_queue_count(self):
"acquire_rotation_index_test",
queues=["default", "compute", "io"],
)
- backend._rotation_offset = 4
+ backend._idle_backoff.rotation_offset = 4
try:
backend.enqueue(replace(echo, queue_name="compute"), args=[1])
backend.enqueue(replace(echo, queue_name="io"), args=[2])
@@ -1185,7 +1190,7 @@ def test_acquire__rebase_rotation_index_beyond_queue_count(self):
assert running_by_queue["compute"] == {acquired[0].id}
assert running_by_queue["io"] == {acquired[1].id}
assert running_by_queue["default"] == set()
- assert backend._rotation_offset == (4 + 2) % 3
+ assert backend._idle_backoff.rotation_offset == (4 + 2) % 3
finally:
backend.close()
@@ -1198,7 +1203,7 @@ def test_acquire__wraps_rotation_offset_through_every_queue(self):
)
recorder = RecordingAcquireScript(backend._acquire_script)
backend._acquire_script = recorder
- backend._rotation_offset = 2
+ backend._idle_backoff.rotation_offset = 2
try:
for repeat in range(2):
for queue_name in queue_names:
@@ -1209,7 +1214,7 @@ def test_acquire__wraps_rotation_offset_through_every_queue(self):
acquired.append(
backend.acquire(*queue_names, timeout=datetime.timedelta(seconds=1))
)
- assert 0 <= backend._rotation_offset < len(queue_names)
+ assert 0 <= backend._idle_backoff.rotation_offset < len(queue_names)
assert [result.task.queue_name for result in acquired] == [
"io",
@@ -1232,7 +1237,7 @@ def test_acquire__keep_rotation_on_idle_polls(self):
)
recorder = RecordingAcquireScript(backend._acquire_script)
backend._acquire_script = recorder
- backend._rotation_offset = 5
+ backend._idle_backoff.rotation_offset = 5
try:
with pytest.raises(TimeoutError):
backend.acquire(
@@ -1242,11 +1247,11 @@ def test_acquire__keep_rotation_on_idle_polls(self):
)
assert len(recorder.sent_args) > 1
assert {sent_args[-1] for sent_args in recorder.sent_args} == {"5"}
- assert backend._rotation_offset == 5
+ assert backend._idle_backoff.rotation_offset == 5
finally:
backend.close()
- def test_init__randomize_rotation_offset(self):
+ def test_idle_backoff__randomize_rotation_offset(self):
"""Seed each backend differently so recycled workers spread across queues."""
queues = ["default", "compute", "io"]
backends = [
@@ -1254,7 +1259,44 @@ def test_init__randomize_rotation_offset(self):
for index in range(32)
]
try:
- assert len({backend._rotation_offset for backend in backends}) > 1
+ assert (
+ len({backend._idle_backoff.rotation_offset for backend in backends}) > 1
+ )
finally:
for backend in backends:
backend.close()
+
+ def test_idle_backoff__isolate_between_threads(self):
+ """Give each worker thread its own idle polling state."""
+ backend = _make_backend("acquire_thread_isolation_test")
+ observed: dict[str, list[IdleBackoff]] = {}
+
+ def observe(name: str) -> None:
+ first = backend._idle_backoff
+ first.miss_count = 7
+ first.rotation_offset = 11
+ observed[name] = [first, backend._idle_backoff]
+
+ thread_names = ("worker-a", "worker-b")
+ threads = [
+ threading.Thread(target=observe, args=(name,)) for name in thread_names
+ ]
+ try:
+ for thread in threads:
+ thread.start()
+ for thread in threads:
+ thread.join()
+
+ # Repeated access on one thread returns the same mutated object.
+ assert observed["worker-a"][0] is observed["worker-a"][1]
+ assert observed["worker-b"][0] is observed["worker-b"][1]
+ assert observed["worker-a"][0].miss_count == 7
+
+ # Separate threads never share state.
+ assert observed["worker-a"][0] is not observed["worker-b"][0]
+
+ # The acquiring thread stays isolated from every worker thread.
+ assert backend._idle_backoff.miss_count == 0
+ assert backend._idle_backoff is not observed["worker-a"][0]
+ finally:
+ backend.close()
diff --git a/tests/test_executor.py b/tests/test_executor.py
index 19ecbbf..fb2e95d 100644
--- a/tests/test_executor.py
+++ b/tests/test_executor.py
@@ -4,6 +4,7 @@
import logging
import multiprocessing
import sys
+import sysconfig
import threading
import time
import uuid
@@ -27,6 +28,7 @@
count_users,
echo,
log_message,
+ record_execution,
)
from threadmill.backends.base import Broker
from threadmill.executor import (
@@ -36,6 +38,9 @@
WorkerThread,
configure_logging,
handler,
+ is_free_threaded_build,
+ is_gil_enabled,
+ warn_when_free_threading_is_unavailable,
)
@@ -198,6 +203,67 @@ def test_configure_logging__keeps_placeholder_loggers(self):
)
+class TestFreeThreadingDiagnostic:
+ """Tests for the interpreter free-threading diagnostic."""
+
+ def test_is_free_threaded_build__report_interpreter(self, monkeypatch):
+ """Report the free-threading build flag of the interpreter."""
+ monkeypatch.setattr(sysconfig, "get_config_var", lambda name: 1)
+ assert is_free_threaded_build() is True
+ monkeypatch.setattr(sysconfig, "get_config_var", lambda name: 0)
+ assert is_free_threaded_build() is False
+
+ def test_is_free_threaded_build__treat_absent_flag_as_gil_build(self, monkeypatch):
+ """Treat an interpreter without the build flag as a GIL build."""
+ monkeypatch.setattr(sysconfig, "get_config_var", lambda name: None)
+ assert is_free_threaded_build() is False
+
+ def test_is_gil_enabled__report_interpreter(self):
+ """Report the GIL state of the running interpreter."""
+ assert is_gil_enabled() is sys._is_gil_enabled()
+
+ def test_is_gil_enabled__assume_enabled_without_private_api(self, monkeypatch):
+ """Assume the GIL is enabled when the private API is absent."""
+ monkeypatch.delattr(sys, "_is_gil_enabled", raising=False)
+ assert is_gil_enabled() is True
+
+ def test_warn_when_free_threading_is_unavailable__warn_lost_parallelism(
+ self, monkeypatch, caplog
+ ):
+ """Warn when a free-threaded build runs with the GIL enabled."""
+ monkeypatch.setattr("threadmill.executor.is_free_threaded_build", lambda: True)
+ monkeypatch.setattr("threadmill.executor.is_gil_enabled", lambda: True)
+
+ with caplog.at_level(logging.WARNING, logger="multiprocessing"):
+ warn_when_free_threading_is_unavailable()
+
+ assert "so worker threads run one at a time" in caplog.text
+
+ def test_warn_when_free_threading_is_unavailable__silent_without_gil(
+ self, monkeypatch, caplog
+ ):
+ """Stay silent when a free-threaded build runs without the GIL."""
+ monkeypatch.setattr("threadmill.executor.is_free_threaded_build", lambda: True)
+ monkeypatch.setattr("threadmill.executor.is_gil_enabled", lambda: False)
+
+ with caplog.at_level(logging.WARNING, logger="multiprocessing"):
+ warn_when_free_threading_is_unavailable()
+
+ assert "so worker threads run one at a time" not in caplog.text
+
+ def test_warn_when_free_threading_is_unavailable__silent_on_gil_build(
+ self, monkeypatch, caplog
+ ):
+ """Stay silent on an interpreter without free threading."""
+ monkeypatch.setattr("threadmill.executor.is_free_threaded_build", lambda: False)
+ monkeypatch.setattr("threadmill.executor.is_gil_enabled", lambda: True)
+
+ with caplog.at_level(logging.WARNING, logger="multiprocessing"):
+ warn_when_free_threading_is_unavailable()
+
+ assert "so worker threads run one at a time" not in caplog.text
+
+
class TestTaskExecutor:
"""Tests for the TaskExecutor dataclass and its methods."""
@@ -287,6 +353,32 @@ def test_run__processes_enqueued_tasks_end_to_end(self):
assert {r.id for r in results} == {r.id for r in enqueued}
assert all(r.status == TaskResultStatus.SUCCESSFUL for r in results)
+ def test_run__processes_each_task_exactly_once_with_threads(self, tmp_path):
+ """Process every queued task exactly once across concurrent worker threads."""
+ count = 40
+ execution_log = tmp_path / "executed"
+ for value in range(count):
+ default_task_backend.enqueue(
+ record_execution, args=[str(execution_log), value]
+ )
+
+ executor = TaskExecutor(
+ backend=default_task_backend,
+ workers=1,
+ threads=4,
+ queues=("default",),
+ exit_empty=True,
+ )
+ run_thread = threading.Thread(target=executor.run, daemon=True)
+ run_thread.start()
+ run_thread.join(timeout=60)
+ assert not run_thread.is_alive()
+
+ recorded = sorted(
+ int(line) for line in execution_log.read_text(encoding="utf-8").split()
+ )
+ assert recorded == list(range(count))
+
def test_run__routes_task_logs_to_stdout(self, capfd):
"""Emit task log records as JSON on standard output."""
enqueued = default_task_backend.enqueue(log_message, args=["hello from task"])
@@ -576,6 +668,7 @@ def test_retry_delay__none_when_callback_raises(self, caplog) -> None:
"""Return None and log when the retry callback raises an exception."""
mp_logger = multiprocessing.get_logger()
mp_logger.addHandler(caplog.handler)
+ original_level = mp_logger.level
mp_logger.setLevel(logging.ERROR)
result = _task_result(boom_retry_raises)
result = dataclasses.replace(
@@ -587,6 +680,7 @@ def test_retry_delay__none_when_callback_raises(self, caplog) -> None:
assert WorkerThread.retry_delay(result) is None
finally:
mp_logger.removeHandler(caplog.handler)
+ mp_logger.setLevel(original_level)
assert "Retry callback failed" in caplog.text
def test_retry_delay__passes_task_context(self) -> None:
diff --git a/tests/testapp/tasks.py b/tests/testapp/tasks.py
index e8481cb..f1d71e4 100644
--- a/tests/testapp/tasks.py
+++ b/tests/testapp/tasks.py
@@ -29,6 +29,13 @@ def log_message(message):
return message
+@task()
+def record_execution(path, value):
+ """Append one line per execution so tests can detect duplicate processing."""
+ with open(path, "a", encoding="utf-8") as recorded:
+ recorded.write(f"{value}\n")
+
+
@task()
def count_users():
"""Count all users in the database (tests model access in workers)."""
diff --git a/threadmill/backends/redis.py b/threadmill/backends/redis.py
index 9ab1a80..0cb0401 100644
--- a/threadmill/backends/redis.py
+++ b/threadmill/backends/redis.py
@@ -6,6 +6,7 @@
import logging
import queue
import random
+import threading
import time
import uuid
from collections.abc import Generator, Sequence
@@ -39,6 +40,18 @@ def _load_lua(name: str) -> str:
return (_LUA_DIR / f"{name}.lua").read_text()
+@dataclasses.dataclass(kw_only=True, slots=True)
+class IdleBackoff:
+ """Idle polling state for a single worker thread.
+
+ Every worker thread backs off and rotates independently, so one thread's
+ misses never reset or double another thread's delay.
+ """
+
+ miss_count: int = 0
+ rotation_offset: int = 0
+
+
class RedisBroker(Broker):
"""Background maintenance broker for the Redis backend."""
@@ -162,11 +175,21 @@ def __init__(self, alias: str, params: dict) -> None:
self.poll_max_interval = self.options.get(
"poll_max_interval", datetime.timedelta(seconds=1)
)
- self._miss_count = 0
- self._rotation_offset = random.randrange(len(self.queues)) # noqa: S311
+ self._idle_backoffs = threading.local()
self._acquire_script = self.client.register_script(self.ACQUIRE_SCRIPT)
self._acknowledge_script = self.client.register_script(self.ACKNOWLEDGE_SCRIPT)
+ @property
+ def _idle_backoff(self) -> IdleBackoff:
+ """Return the idle polling state of the calling worker thread."""
+ backoff = getattr(self._idle_backoffs, "backoff", None)
+ if backoff is None:
+ backoff = IdleBackoff(
+ rotation_offset=random.randrange(len(self.queues)) # noqa: S311
+ )
+ self._idle_backoffs.backoff = backoff
+ return backoff
+
@property
def async_client(self) -> redis.asyncio.Redis:
"""Lazily-created async Redis client, reused across calls."""
@@ -263,6 +286,7 @@ def acquire(
)
]
+ idle_backoff = self._idle_backoff
while True:
now = timezone.now()
now_ms = now.timestamp() * 1000
@@ -277,11 +301,13 @@ def acquire(
str(len(queue_names)),
worker,
str(int(self.lease_ttl.total_seconds() * 1000)),
- str(self._rotation_offset),
+ str(idle_backoff.rotation_offset),
],
):
- self._miss_count = 0
- self._rotation_offset = (self._rotation_offset + 1) % len(queue_names)
+ idle_backoff.miss_count = 0
+ idle_backoff.rotation_offset = (idle_backoff.rotation_offset + 1) % len(
+ queue_names
+ )
return self.deserialize_task_result(data)
try:
@@ -294,11 +320,12 @@ def acquire(
# exponents would only overflow the float math.
cap = int(self.poll_max_interval / self.poll_interval).bit_length()
interval_secs = min(
- self.poll_interval.total_seconds() * 2 ** min(self._miss_count, cap),
+ self.poll_interval.total_seconds()
+ * 2 ** min(idle_backoff.miss_count, cap),
self.poll_max_interval.total_seconds(),
remaining,
)
- self._miss_count += 1
+ idle_backoff.miss_count += 1
time.sleep(interval_secs)
def acknowledge(self, task_result: TaskResult) -> None:
diff --git a/threadmill/executor.py b/threadmill/executor.py
index 1f4a58c..5d33446 100644
--- a/threadmill/executor.py
+++ b/threadmill/executor.py
@@ -9,6 +9,7 @@
import random
import socket
import sys
+import sysconfig
import threading
import time
import typing
@@ -99,6 +100,31 @@ def configure_logging(formatter: logging.Formatter) -> None:
root_logger.setLevel(logging.INFO)
+def is_free_threaded_build() -> bool:
+ """Return whether this interpreter is a free-threaded build."""
+ return bool(sysconfig.get_config_var("Py_GIL_DISABLED"))
+
+
+def is_gil_enabled() -> bool:
+ """Return whether the global interpreter lock is currently enabled."""
+ return getattr(sys, "_is_gil_enabled", lambda: True)()
+
+
+def warn_when_free_threading_is_unavailable() -> None:
+ """Warn when a free-threaded build runs with the GIL enabled.
+
+ A C extension that has not declared free-threading support re-enables the GIL
+ for the whole process, which silently costs the parallelism the build exists
+ to provide.
+ """
+ if is_free_threaded_build() and is_gil_enabled():
+ logger.warning(
+ "free-threaded interpreter is running with the GIL enabled, so worker "
+ "threads run one at a time. A C extension that does not declare "
+ "free-threading support, such as hiredis, re-enables the GIL"
+ )
+
+
@dataclasses.dataclass(kw_only=True, slots=True)
class TaskExecutor:
"""Tasks consumed from shared joinable queues via process and thread pools."""
@@ -151,6 +177,9 @@ def create_worker_process(self) -> WorkerProcess:
def run(self) -> None:
"""Start consuming tasks until shutdown is requested."""
configure_logging(self.log_formatter)
+ # The warning only reaches the log handler once configure_logging has
+ # routed the multiprocessing logger through the root logger.
+ warn_when_free_threading_is_unavailable()
self.worker_processes = [
self.create_worker_process() for _ in range(self.process_count)
]
From fd05f89665bcd09d28cd2ed687b2839de8ca86c0 Mon Sep 17 00:00:00 2001
From: Johannes Maron
Date: Thu, 1 Oct 2026 16:17:24 +0200
Subject: [PATCH 2/7] Document the free-threading trade-off and keep the worker
defaults
The scaling benchmark shows that one process with four threads reaches
the same throughput on a free-threaded build as four processes with one
thread each, so the existing defaults already reach full parallelism on
3.14t. Auto-tuning them would buy a smaller memory and connection
footprint rather than throughput, at the cost of confining a crash or a
task recycling to a single process.
Record that reasoning next to the numbers that produced it, and tell
users which way to move and what they give up when they do.
---
README.md | 2 ++
benchmarks/test_scaling.py | 8 ++++++++
2 files changed, 10 insertions(+)
diff --git a/README.md b/README.md
index 6f99bc0..65dcbb0 100644
--- a/README.md
+++ b/README.md
@@ -89,6 +89,8 @@ A pool on a free-threaded interpreter therefore reaches the same throughput with
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]
diff --git a/benchmarks/test_scaling.py b/benchmarks/test_scaling.py
index 9ae207d..635b3eb 100644
--- a/benchmarks/test_scaling.py
+++ b/benchmarks/test_scaling.py
@@ -16,6 +16,14 @@
includes the worker pool's fixed start cost of one to two seconds, which is a
larger share of the faster configurations, so compare each configuration against
the ``(1, 1)`` baseline instead of reading the times as pure throughput.
+
+Measured on 16 CPU-bound tasks, one process with four threads reaches the same
+throughput as four processes with one thread each, while four of each reaches
+4.2x. A GIL build gains nothing from extra threads at all. That is why the worker
+defaults stay one process per core with a single thread: they already reach full
+parallelism on a free-threaded build, and a crash or a task recycling stays
+confined to one process. Threads trade that isolation for a smaller memory and
+connection footprint, which is a per-deployment choice rather than a default.
"""
import dataclasses
From 84d1709e539e28346e53542290158e57f872dafa Mon Sep 17 00:00:00 2001
From: Johannes Maron
Date: Thu, 8 Oct 2026 00:19:51 +0200
Subject: [PATCH 3/7] Benchmark CPU parallelism against other frameworks
The queue comparison measures queue overhead with a trivial echo task, and
the thread scaling benchmark measures threadmill alone, so neither answers
whether threadmill's free-threading parallelism is unusual.
Runs the same CPU-bound workload through every framework that can spread
work across threads, one thread and four on both interpreters, 16 tasks of
about a second each. Measured on this machine:
build framework 1 thread 4 threads ratio
3.14 (GIL) threadmill 17.08s 17.07s 1.00x
3.14 (GIL) celery 15.85s 15.38s 1.03x
3.14 (GIL) dramatiq 15.25s 15.24s 1.00x
3.14t threadmill 17.15s 6.07s 2.83x
3.14t celery 15.85s 4.87s 3.26x
3.14t dramatiq 14.99s 4.62s 3.25x
No framework gains anything from threads on a GIL build, which is the
control that says the measurement is sound. All of them gain real
parallelism on a free-threaded build.
threadmill's lower ratio is its start cost, not weaker parallelism. Its
process pool costs about two seconds to start, against a fraction of a
second for the other two, and that cost sits in both drains. Subtracting
it, four threads run the work 3.7x faster, against 3.5x for dramatiq and
3.8x for celery. Its four-thread and four-process drains are also equal,
6.07s against 6.08s, so threads substitute for processes.
Three fixes the measurements forced:
- The workload drains its own queue on every framework. A worker process
left behind by another worktree sharing the same Redis instance consumed
the Celery queue and made a drain wait for a task it had already run.
- Completion is counted by the tasks themselves rather than by a sentinel
consumed last. With several threads a free thread takes the sentinel
before its predecessors finish, which would report a drain that never
did the work.
- One drain per configuration. The tasks are queued once, so a second
round finds an empty queue and measures only the start cost, and a
median over rounds averages work against no work.
---
benchmarks/celery_app.py | 9 +
benchmarks/cpu_work.py | 52 ++++
benchmarks/dramatiq_app.py | 9 +
benchmarks/test_parallelism.py | 465 +++++++++++++++++++++++++++++++++
4 files changed, 535 insertions(+)
create mode 100644 benchmarks/cpu_work.py
create mode 100644 benchmarks/test_parallelism.py
diff --git a/benchmarks/celery_app.py b/benchmarks/celery_app.py
index 49c67a9..3f3eb27 100644
--- a/benchmarks/celery_app.py
+++ b/benchmarks/celery_app.py
@@ -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"
@@ -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."""
diff --git a/benchmarks/cpu_work.py b/benchmarks/cpu_work.py
new file mode 100644
index 0000000..b68ad38
--- /dev/null
+++ b/benchmarks/cpu_work.py
@@ -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
diff --git a/benchmarks/dramatiq_app.py b/benchmarks/dramatiq_app.py
index b3f474a..5833f32 100644
--- a/benchmarks/dramatiq_app.py
+++ b/benchmarks/dramatiq_app.py
@@ -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"
@@ -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)
diff --git a/benchmarks/test_parallelism.py b/benchmarks/test_parallelism.py
new file mode 100644
index 0000000..0a8cd17
--- /dev/null
+++ b/benchmarks/test_parallelism.py
@@ -0,0 +1,465 @@
+"""Benchmark how much CPU parallelism each framework delivers.
+
+The queue comparison in ``test_backends.py`` measures queue overhead with a
+trivial echo task, and ``test_scaling.py`` measures threadmill alone. Neither
+answers whether threadmill's free-threading parallelism is unusual, so this
+module runs the same CPU-bound workload through every framework that can spread
+work across threads and reports how long each takes.
+
+Only a thread-based configuration can use a free-threaded interpreter:
+
+- threadmill: one process with N threads
+- dramatiq: one process with N threads, its native model
+- celery: ``--pool=threads --concurrency=N``
+
+django-tasks-db processes one task at a time, and RQ forks a process for each
+job, so neither has a thread-based configuration to compare.
+
+Every configuration is measured with one thread and with four, from the same
+framework and the same queue. The ratio between the two is the useful comparison
+rather than the raw times, because each framework carries a different fixed start
+cost and that cost sits in both drains. threadmill starts a pool of processes and
+pays for the forkserver, the interpreter and Django setup in each one, which costs
+seconds, while Celery and dramatiq start a worker in a fraction of a second. A
+framework with a cheaper start therefore shows a larger ratio at the same
+parallelism. Read the ratio as the speedup of the whole pool, and subtract the
+start cost before reading it as the parallelism of the work.
+
+Each drain is measured until every task reported completion, not until a sentinel
+was seen: with several threads a free thread can take a sentinel before its
+predecessors finish, which would report a drain that never did the work.
+
+Run both interpreters and compare:
+
+ uv run pytest benchmarks/test_parallelism.py -m benchmark
+ uv run --python 3.14t pytest benchmarks/test_parallelism.py -m benchmark
+"""
+
+import collections.abc
+import dataclasses
+import os
+import subprocess
+import sys
+import tempfile
+import time
+import typing
+
+import pytest
+import redis
+from django.core.management import call_command
+from django.tasks import (
+ DEFAULT_TASK_BACKEND_ALIAS,
+ TaskResultStatus,
+ task_backends,
+)
+
+from benchmarks import cpu_work
+from benchmarks.celery_app import REDIS_URL, celery_compute
+from benchmarks.dramatiq_app import dramatiq_compute
+from tests.testapp.tasks import compute_workload
+
+TASK_COUNT = 16
+"""CPU-bound tasks drained per measurement, about one second of work each.
+
+Deep enough that every worker in every configuration gets a share. threadmill
+reads ahead by ``threads * 4``, so one process running four threads asks for the
+whole queue and spreads it over its threads, while four processes running one
+thread each ask for four apiece and spread it evenly over the processes. With
+fewer tasks the process pool is lopsided, because whichever process starts first
+takes the queue and the rest boot into an empty one.
+"""
+
+THREAD_COUNT = 4
+"""Threads in the parallel configuration, and the speedup denominator is one."""
+
+MEASUREMENT_ROUNDS = 1
+"""Drains per configuration.
+
+Each drain starts a pool or a worker process and processes the queue, and the
+tasks are queued once. A second round would find an empty queue and only pay the
+start cost again, so a median over rounds would average work against no work. The
+external frameworks also stop their worker outside the measured region, so a
+second round would measure two workers competing.
+"""
+
+SCALING_QUEUE_NAME = "scaling"
+"""Queue dedicated to the CPU benchmarks.
+
+A separate queue keeps the measurement honest: a worker process left behind by
+another test or worktree sharing the same Redis instance would otherwise steal
+tasks and report a drain that never did the work.
+"""
+
+WORKER_STOP_TIMEOUT_SECONDS = 20
+"""Seconds to wait for a worker process to stop after SIGTERM."""
+
+DRAIN_TIMEOUT_SECONDS = 300
+"""Seconds to wait for a worker to process every queued task."""
+
+READ_AHEAD = 128
+"""Messages each worker reads ahead, so a drain pulls its whole queue at once."""
+
+CELERY_SINGLE_THREAD_WORKER = (
+ sys.executable,
+ "-m",
+ "celery",
+ "-A",
+ "benchmarks.celery_app:celery_app",
+ "worker",
+ "--pool=solo",
+ "-Q",
+ cpu_work.CPU_QUEUE,
+ f"--prefetch-multiplier={READ_AHEAD}",
+ "--loglevel=WARNING",
+ "--without-gossip",
+ "--without-mingle",
+ "--without-heartbeat",
+)
+"""Celery running one task at a time. The prefork pool crashes on CPython 3.14."""
+
+CELERY_MULTI_THREAD_WORKER = (
+ sys.executable,
+ "-m",
+ "celery",
+ "-A",
+ "benchmarks.celery_app:celery_app",
+ "worker",
+ "--pool=threads",
+ f"--concurrency={THREAD_COUNT}",
+ "-Q",
+ cpu_work.CPU_QUEUE,
+ f"--prefetch-multiplier={READ_AHEAD}",
+ "--loglevel=WARNING",
+ "--without-gossip",
+ "--without-mingle",
+ "--without-heartbeat",
+)
+"""Celery running its thread pool."""
+
+DRAMATIQ_SINGLE_THREAD_WORKER = (
+ sys.executable,
+ "-m",
+ "dramatiq",
+ "benchmarks.dramatiq_app:redis_broker",
+ "--processes",
+ "1",
+ "--threads",
+ "1",
+ "--queues",
+ cpu_work.CPU_QUEUE,
+)
+"""dramatiq running one process with one thread, its smallest configuration."""
+
+DRAMATIQ_MULTI_THREAD_WORKER = (
+ sys.executable,
+ "-m",
+ "dramatiq",
+ "benchmarks.dramatiq_app:redis_broker",
+ "--processes",
+ "1",
+ "--threads",
+ str(THREAD_COUNT),
+ "--queues",
+ cpu_work.CPU_QUEUE,
+)
+"""dramatiq running one process with several threads, its native model."""
+
+DRAMATIQ_ENV = {**os.environ, "dramatiq_queue_prefetch": str(READ_AHEAD)}
+"""Environment carrying the dramatiq read-ahead, which has no command line flag."""
+
+scaling_workload = dataclasses.replace(compute_workload, queue_name=SCALING_QUEUE_NAME)
+"""The CPU-bound Django task, pinned to the benchmark's own queue."""
+
+
+@dataclasses.dataclass(frozen=True, kw_only=True, slots=True)
+class WorkerProcess:
+ """A worker subprocess started by a benchmark, with its captured log."""
+
+ process: subprocess.Popen
+ log: typing.IO[bytes]
+
+
+running_workers: list[WorkerProcess] = []
+"""Workers this benchmark started, stopped once each measurement ends."""
+
+
+def enqueue_threadmill(count: int) -> None:
+ """Queue `count` CPU tasks on the threadmill backend."""
+ queued_task = scaling_workload.using(backend=DEFAULT_TASK_BACKEND_ALIAS)
+ for _ in range(count):
+ queued_task.enqueue()
+
+
+def enqueue_celery(count: int) -> None:
+ """Queue `count` CPU tasks on the Celery broker."""
+ for index in range(count):
+ celery_compute.delay(index)
+
+
+def enqueue_dramatiq(count: int) -> None:
+ """Queue `count` CPU tasks on the dramatiq broker."""
+ for index in range(count):
+ dramatiq_compute.send(index)
+
+
+def drain_with_threadmill_worker() -> None:
+ """Drain the queue with one threadmill process running one thread."""
+ _drain_with_threadmill_pool(workers=1, threads=1)
+
+
+def drain_with_threadmill_threads() -> None:
+ """Drain the queue with one threadmill process running several threads."""
+ _drain_with_threadmill_pool(workers=1, threads=THREAD_COUNT)
+
+
+def drain_with_threadmill_processes() -> None:
+ """Drain the queue with several threadmill processes running one thread each."""
+ _drain_with_threadmill_pool(workers=THREAD_COUNT, threads=1)
+
+
+def _drain_with_threadmill_pool(*, workers: int, threads: int) -> None:
+ """Drain the queue and exit, so the measurement includes the pool start cost."""
+ call_command(
+ "threadmill",
+ "worker",
+ backend=DEFAULT_TASK_BACKEND_ALIAS,
+ queues=[SCALING_QUEUE_NAME],
+ workers=workers,
+ threads=threads,
+ exit_empty=True,
+ verbosity=0,
+ )
+
+
+def drain_with_celery_single_thread() -> None:
+ """Drain the queue with a one-thread Celery worker."""
+ drain_with_external_worker(CELERY_SINGLE_THREAD_WORKER)
+
+
+def drain_with_celery_threads() -> None:
+ """Drain the queue with a multi-thread Celery worker."""
+ drain_with_external_worker(CELERY_MULTI_THREAD_WORKER)
+
+
+def drain_with_dramatiq_single_thread() -> None:
+ """Drain the queue with a one-thread dramatiq worker."""
+ drain_with_external_worker(DRAMATIQ_SINGLE_THREAD_WORKER, env=DRAMATIQ_ENV)
+
+
+def drain_with_dramatiq_threads() -> None:
+ """Drain the queue with a multi-thread dramatiq worker."""
+ drain_with_external_worker(DRAMATIQ_MULTI_THREAD_WORKER, env=DRAMATIQ_ENV)
+
+
+def drain_with_external_worker(
+ argv: collections.abc.Sequence[str],
+ env: collections.abc.Mapping[str, str] | None = None,
+) -> None:
+ """Run a worker CLI until every queued task reported completion.
+
+ The worker is stopped after the measurement, by the ``stop_workers`` fixture,
+ because a graceful shutdown takes seconds and would dominate a short drain.
+ """
+ client = redis.Redis.from_url(REDIS_URL)
+ client.delete(cpu_work.COMPLETION_KEY)
+ log = tempfile.TemporaryFile()
+ # The command is a fixed worker CLI, never caller input.
+ process = subprocess.Popen( # noqa: S603
+ argv,
+ env=env,
+ stdout=log,
+ stderr=subprocess.STDOUT,
+ )
+ running_workers.append(WorkerProcess(process=process, log=log))
+ wait_until_completed(client, process, log)
+
+
+def wait_until_completed(
+ client: redis.Redis, process: subprocess.Popen, log: typing.IO[bytes]
+) -> None:
+ """Wait until every CPU task incremented the completion counter."""
+ deadline = time.monotonic() + DRAIN_TIMEOUT_SECONDS
+ while time.monotonic() < deadline:
+ if int(client.get(cpu_work.COMPLETION_KEY) or 0) >= TASK_COUNT:
+ return
+ if process.poll() is not None:
+ raise AssertionError(
+ f"Worker exited with {process.returncode} after "
+ f"{int(client.get(cpu_work.COMPLETION_KEY) or 0)} of {TASK_COUNT}"
+ f" tasks:\n{read_log_tail(log)}"
+ )
+ time.sleep(0.001)
+ raise AssertionError(
+ f"Worker processed {int(client.get(cpu_work.COMPLETION_KEY) or 0)} of"
+ f" {TASK_COUNT} tasks within {DRAIN_TIMEOUT_SECONDS}s:\n{read_log_tail(log)}"
+ )
+
+
+def read_log_tail(log: typing.IO[bytes], line_count: int = 40) -> str:
+ """Return the last lines of a worker log."""
+ log.seek(0)
+ lines = log.read().decode(errors="replace").splitlines()
+ return "\n".join(lines[-line_count:])
+
+
+def stop_process(process: subprocess.Popen) -> None:
+ """Ask a worker to stop, killing it if it does not exit in time."""
+ if process.poll() is not None:
+ return
+ process.terminate()
+ try:
+ process.wait(timeout=WORKER_STOP_TIMEOUT_SECONDS)
+ except subprocess.TimeoutExpired:
+ process.kill()
+ process.wait(timeout=WORKER_STOP_TIMEOUT_SECONDS)
+
+
+def verify_threadmill_drain() -> None:
+ """Assert the threadmill pool left no task queued and every task succeeded."""
+ backend = task_backends[DEFAULT_TASK_BACKEND_ALIAS]
+ for status in (TaskResultStatus.READY, TaskResultStatus.RUNNING):
+ remaining = sum(
+ 1 for _ in backend.peek(SCALING_QUEUE_NAME, status=status, count=0)
+ )
+ assert remaining == 0, f"{remaining} tasks left in {status.name}"
+ successful = sum(
+ 1
+ for _ in backend.peek(
+ SCALING_QUEUE_NAME, status=TaskResultStatus.SUCCESSFUL, count=0
+ )
+ )
+ assert successful == TASK_COUNT, f"{successful} of {TASK_COUNT} tasks succeeded"
+
+
+@dataclasses.dataclass(frozen=True, kw_only=True, slots=True)
+class FrameworkParallelism:
+ """One framework and thread configuration to measure."""
+
+ name: str
+ """Identifier in the benchmark report."""
+
+ enqueue: collections.abc.Callable[[int], None]
+ """Queue the CPU workload."""
+
+ drain: collections.abc.Callable[[], None]
+ """Process the queued workload."""
+
+ verify: collections.abc.Callable[[], None] | None = None
+ """Assert the drain did the work, where the drain does not already prove it."""
+
+
+PARALLELISMS_UNDER_TEST = (
+ FrameworkParallelism(
+ name="threadmill-1t",
+ enqueue=enqueue_threadmill,
+ drain=drain_with_threadmill_worker,
+ verify=verify_threadmill_drain,
+ ),
+ FrameworkParallelism(
+ name="threadmill-4t",
+ enqueue=enqueue_threadmill,
+ drain=drain_with_threadmill_threads,
+ verify=verify_threadmill_drain,
+ ),
+ FrameworkParallelism(
+ name="threadmill-4p",
+ enqueue=enqueue_threadmill,
+ drain=drain_with_threadmill_processes,
+ verify=verify_threadmill_drain,
+ ),
+ FrameworkParallelism(
+ name="celery-1t", enqueue=enqueue_celery, drain=drain_with_celery_single_thread
+ ),
+ FrameworkParallelism(
+ name="celery-4t", enqueue=enqueue_celery, drain=drain_with_celery_threads
+ ),
+ FrameworkParallelism(
+ name="dramatiq-1t",
+ enqueue=enqueue_dramatiq,
+ drain=drain_with_dramatiq_single_thread,
+ ),
+ FrameworkParallelism(
+ name="dramatiq-4t",
+ enqueue=enqueue_dramatiq,
+ drain=drain_with_dramatiq_threads,
+ ),
+)
+"""Configurations spanning each framework at one thread and at several.
+
+The queue overhead benchmark already compares threadmill, django-tasks-db and
+django-tasks-rq for throughput with an echo task, so this one covers the
+thread-based configurations those two cannot offer.
+"""
+
+
+def identify_parallelism(parallelism: FrameworkParallelism) -> str:
+ """Return the benchmark identifier of a configuration."""
+ return parallelism.name
+
+
+@pytest.fixture(autouse=True)
+def stop_workers(empty_queues):
+ """Stop every worker a measurement started, after the measurement.
+
+ Depends on ``empty_queues`` so that the cleanup runs after the workers stop,
+ rather than while one is still writing to Redis.
+ """
+ yield
+ while running_workers:
+ worker = running_workers.pop()
+ stop_process(worker.process)
+ worker.log.close()
+
+
+@pytest.fixture
+def empty_queues():
+ """Delete queued tasks from every compared queue before and after a measurement."""
+ client = redis.Redis.from_url(REDIS_URL)
+ threadmill_client = task_backends[DEFAULT_TASK_BACKEND_ALIAS].client
+
+ def delete_queued_tasks() -> None:
+ if keys := threadmill_client.keys("threadmill:*"):
+ threadmill_client.delete(*keys)
+ # The Celery queue is a bare list named after the queue, so it has no
+ # prefix to match and is named here.
+ for key_pattern in ("celery*", "_kombu*", cpu_work.CPU_QUEUE, "dramatiq:*"):
+ if keys := client.keys(key_pattern):
+ client.delete(*keys)
+ client.delete(cpu_work.COMPLETION_KEY)
+
+ delete_queued_tasks()
+ yield
+ delete_queued_tasks()
+
+
+class TestCpuParallelism:
+ """Measure how long a CPU-bound workload takes to drain."""
+
+ @pytest.mark.benchmark
+ @pytest.mark.django_db(transaction=True)
+ @pytest.mark.parametrize(
+ "parallelism",
+ PARALLELISMS_UNDER_TEST,
+ ids=identify_parallelism,
+ )
+ def test_drain_cpu_workload__benchmark(self, benchmark, parallelism, empty_queues):
+ """Benchmark the time for every task to report completion."""
+ benchmark.extra_info.update(
+ {
+ "tasks": TASK_COUNT,
+ "threads": THREAD_COUNT,
+ "python": sys.version.split()[0],
+ }
+ )
+ parallelism.enqueue(TASK_COUNT)
+
+ benchmark.pedantic(
+ parallelism.drain,
+ rounds=MEASUREMENT_ROUNDS,
+ iterations=1,
+ warmup_rounds=0,
+ )
+
+ # Outside the timed region: a drain that skipped work must not look fast.
+ if parallelism.verify:
+ parallelism.verify()
From 3ed797e1a4b5b00a82426864d121be4f44fb09a1 Mon Sep 17 00:00:00 2001
From: Johannes Maron
Date: Thu, 8 Oct 2026 01:25:50 +0200
Subject: [PATCH 4/7] Chart Threadmill free threading in the README comparison
The chart showed one Threadmill bar, so the headline feature was invisible
next to the other queues.
- Add a Threadmill (free threading) bar, measured with four threads in one
process: 59,697 tasks per second against 11,963 for one thread.
- Capitalise the label as Threadmill.
- List the bar only where the running interpreter really is free-threaded
with the GIL off. Elsewhere the drain would only repeat the
single-threaded rate, which would read as a feature that found no
speedup.
- Stack the footnote over two lines and note the four-thread configuration.
The bar is a free-threading result rather than four threads overlapping IO.
On a GIL build four threads reach 9,970 tasks per second against 7,497 for
one, a gain of 1.33x, while on a free-threaded build the same four threads
reach 58,306 against 11,881, a gain of 4.91x.
---
README.md | 2 +-
benchmarks/chart.py | 62 +++++++++++++++------
benchmarks/test_backends.py | 68 +++++++++++++++++++++++-
docs/images/backend-comparison-dark.svg | 40 +++++++-------
docs/images/backend-comparison-light.svg | 40 +++++++-------
5 files changed, 157 insertions(+), 55 deletions(-)
diff --git a/README.md b/README.md
index 371eda9..2a01993 100644
--- a/README.md
+++ b/README.md
@@ -19,7 +19,7 @@
-
+
diff --git a/benchmarks/chart.py b/benchmarks/chart.py
index f89ec8f..247be90 100644
--- a/benchmarks/chart.py
+++ b/benchmarks/chart.py
@@ -77,7 +77,7 @@ class Theme:
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
@@ -85,6 +85,12 @@ class Theme:
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."""
+
+FOOTNOTE_LINE_HEIGHT = 17
+"""Distance between the baselines of two stacked footnote lines."""
+
@dataclasses.dataclass(frozen=True, kw_only=True, slots=True)
class QueueResult:
@@ -144,11 +150,32 @@ def text(x, y, content, *, theme, size=13, fill=None, weight=400, anchor="start"
def describe(results: list[QueueResult]) -> str:
"""Return a sentence describing the throughput of every queue."""
- return "Tasks per second with one worker: " + ", ".join(
+ return "Tasks per second with one worker process each: " + ", ".join(
f"{result.name} {result.throughput:,.0f}" for result in results
)
+def create_footnote_lines(results: list[QueueResult]) -> list[str]:
+ """Return the footnote, one line per element.
+
+ The worker configuration line names the free-threading queue only when the
+ benchmark measured one, so a chart built from a GIL build claims no thread
+ parallelism it did not measure.
+ """
+ if any(result.name == FREE_THREADING_QUEUE for result in results):
+ workers = (
+ "One process and one thread each, except Threadmill on a free-threaded"
+ " build, which runs four."
+ )
+ else:
+ workers = "One process and one thread each."
+ return [
+ workers,
+ "Threadmill, celery and dramatiq read 128 ahead;"
+ " django-tasks-db and -rq read one message at a time.",
+ ]
+
+
def build_chart(results: list[QueueResult], theme: Theme) -> str:
"""Return the chart as an SVG document drawn in the given theme."""
fastest = results[0]
@@ -177,8 +204,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,
@@ -217,19 +244,22 @@ def build_chart(results: list[QueueResult], theme: Theme) -> str:
)
)
- parts.append(
- text(
- 28,
- footnote_y,
- # 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,
+ footnote_lines = create_footnote_lines(results)
+ for line_index, line in enumerate(footnote_lines):
+ parts.append(
+ text(
+ 28,
+ footnote_y + line_index * FOOTNOTE_LINE_HEIGHT,
+ # joe: widths estimated against the 11.5px hand check this comment
+ # replaced (right edge 856.1 of 900 for 132 characters); the longest
+ # line now reaches about 675 of 900. Recheck if the canvas width or
+ # the font stack changes.
+ line,
+ theme=theme,
+ size=11.5,
+ fill=theme.faint,
+ )
)
- )
parts.append("")
return "\n".join(parts)
diff --git a/benchmarks/test_backends.py b/benchmarks/test_backends.py
index 3fbbb57..6b2ffb3 100644
--- a/benchmarks/test_backends.py
+++ b/benchmarks/test_backends.py
@@ -38,6 +38,11 @@
Its queues are deeper than the others, because its marginal drain is only seconds
long. A shallow queue sits inside the one-second quantization of the fixed cost.
+On a free-threaded interpreter with the GIL disabled, Threadmill is measured once
+more with ``FREE_THREADING_THREAD_COUNT`` threads in one process. That queue is
+listed only there, because on any other interpreter the threads run one at a time
+and the drain would only repeat the single-threaded rate.
+
Threadmill, django-tasks-db and django-tasks-rq run one worker process that drains
a queue and exits. Celery and dramatiq have no such mode, so the benchmark queues a
sentinel task last and waits for it. That wait proves that the queue was drained.
@@ -75,6 +80,7 @@
)
from benchmarks.dramatiq_app import dramatiq_echo, dramatiq_mark_processed
from tests.testapp.tasks import echo
+from threadmill.executor import is_free_threaded_build, is_gil_enabled
ENQUEUE_ITERATIONS = 500
"""Tasks enqueued within one enqueue benchmark round."""
@@ -99,6 +105,15 @@
each job. Neither queue can be told to read ahead.
"""
+FREE_THREADING_THREAD_COUNT = 4
+"""Threads the free-threading threadmill worker runs, where the build allows it.
+
+Deliberately a fixed small number rather than every core, so the bar states what
+a free-threaded build does for the same pool and not what this machine happens to
+have. The other queues all run one thread, which is why the bar is labelled with
+its thread count.
+"""
+
CELERY_WORKER = (
sys.executable,
"-m",
@@ -209,6 +224,27 @@ def drain_with_threadmill_worker_no_prefetch() -> None:
)
+def drain_with_threadmill_free_threading_worker() -> None:
+ """Process every queued task with one threadmill worker on several threads.
+
+ Only meaningful on a free-threaded interpreter, where the threads run at the
+ same time. The queue is only listed under test when the running interpreter is
+ free-threaded, so this drain never reports the single-threaded rate of a GIL
+ build as if it were a parallel one.
+ """
+ call_command(
+ "threadmill",
+ "worker",
+ backend=DEFAULT_TASK_BACKEND_ALIAS,
+ queues=[DEFAULT_TASK_QUEUE_NAME],
+ workers=1,
+ threads=FREE_THREADING_THREAD_COUNT,
+ prefetch_count=READ_AHEAD,
+ exit_empty=True,
+ verbosity=0,
+ )
+
+
def drain_with_django_tasks_db_worker() -> None:
"""Process every queued task with the django-tasks-db worker."""
call_command(
@@ -373,15 +409,43 @@ def django_task_backend(
therefore keeps the prefetch comparison inside that step.
"""
+THREADS_RUN_IN_PARALLEL = is_free_threaded_build() and not is_gil_enabled()
+"""Whether this interpreter runs Python threads at the same time.
+
+A free-threaded build stops doing so as soon as a C extension that has not
+declared free-threading support enables the GIL, so the build flag alone does not
+answer this.
+"""
+
+FREE_THREADING_QUEUES = (
+ (
+ django_task_backend(
+ "Threadmill (free threading)",
+ DEFAULT_TASK_BACKEND_ALIAS,
+ drain_with_threadmill_free_threading_worker,
+ task_count=THREADMILL_TASK_COUNT,
+ ),
+ )
+ if THREADS_RUN_IN_PARALLEL
+ else ()
+)
+"""The free-threading queue, listed only where threads really are parallel.
+
+On any other interpreter the drain behind this queue measures the same
+single-threaded rate as the queue above it, which would read as a free-threading
+result that found no speedup.
+"""
+
WORKER_QUEUES = (
django_task_backend(
- "threadmill",
+ "Threadmill",
DEFAULT_TASK_BACKEND_ALIAS,
drain_with_threadmill_worker,
task_count=THREADMILL_TASK_COUNT,
),
+ *FREE_THREADING_QUEUES,
django_task_backend(
- "threadmill (no prefetch)",
+ "Threadmill (no prefetch)",
DEFAULT_TASK_BACKEND_ALIAS,
drain_with_threadmill_worker_no_prefetch,
task_count=THREADMILL_TASK_COUNT,
diff --git a/docs/images/backend-comparison-dark.svg b/docs/images/backend-comparison-dark.svg
index c0ce3a5..fda2833 100644
--- a/docs/images/backend-comparison-dark.svg
+++ b/docs/images/backend-comparison-dark.svg
@@ -1,22 +1,26 @@
-