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 @@ - Tasks per second with one worker: threadmill 11,977, dramatiq 7,168, celery 2,183, django-tasks-db 2,154, django-tasks-rq 87. + Tasks per second with one worker process each: Threadmill (free threading) 59,697, Threadmill 11,963, dramatiq 7,182, celery 2,397, django-tasks-db 2,103, django-tasks-rq 74.

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 @@ - + - + Queue throughput -5,000–60,000 trivial tasks per queue · one worker process, one thread · higher is better -threadmill +5,000–60,000 trivial tasks per queue · one worker process each · higher is better +Threadmill (free threading) -11,977/s -dramatiq - -7,168/s -celery - -2,183/s -django-tasks-db - -2,154/s -django-tasks-rq - -87/s -One process and one thread each. Threadmill, celery and dramatiq read 128 ahead; django-tasks-db and -rq read one message at a time. +59,697/s +Threadmill + +11,963/s +dramatiq + +7,182/s +celery + +2,397/s +django-tasks-db + +2,103/s +django-tasks-rq + +74/s +One process and one thread each, except Threadmill on a free-threaded build, which runs four. +Threadmill, celery and dramatiq read 128 ahead; django-tasks-db and -rq read one message at a time. diff --git a/docs/images/backend-comparison-light.svg b/docs/images/backend-comparison-light.svg index a69b065..ae2e4ed 100644 --- a/docs/images/backend-comparison-light.svg +++ b/docs/images/backend-comparison-light.svg @@ -1,22 +1,26 @@ - + - + Queue throughput -5,000–60,000 trivial tasks per queue · one worker process, one thread · higher is better -threadmill +5,000–60,000 trivial tasks per queue · one worker process each · higher is better +Threadmill (free threading) -11,977/s -dramatiq - -7,168/s -celery - -2,183/s -django-tasks-db - -2,154/s -django-tasks-rq - -87/s -One process and one thread each. Threadmill, celery and dramatiq read 128 ahead; django-tasks-db and -rq read one message at a time. +59,697/s +Threadmill + +11,963/s +dramatiq + +7,182/s +celery + +2,397/s +django-tasks-db + +2,103/s +django-tasks-rq + +74/s +One process and one thread each, except Threadmill on a free-threaded build, which runs four. +Threadmill, celery and dramatiq read 128 ahead; django-tasks-db and -rq read one message at a time. From 68b187619c291fa2f87e0875ce531e236d115bc7 Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Thu, 8 Oct 2026 17:37:55 +0200 Subject: [PATCH 5/7] Inline the free-threading checks and shorten the warning - Drop the note above the warning call. The order of the two lines says it. - Inline the interpreter checks. They are single expressions, nothing reused them, and the benchmarks had to import them to report the build. - Warn in two sentences instead of three, without naming a C extension. The cost the GIL inflicts is still in the docstring, where it explains the check. - Drive the warning in the tests through the stdlib attributes it reads, so the tests cover the absent private API again. - Drop the comment above the free-threaded CI job. --- .github/workflows/ci.yml | 1 - benchmarks/test_scaling.py | 8 +++--- tests/test_executor.py | 55 +++++++++++++++----------------------- threadmill/executor.py | 21 ++++----------- 4 files changed, 32 insertions(+), 53 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 77fec46..cda143e 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -67,7 +67,6 @@ jobs: - 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 diff --git a/benchmarks/test_scaling.py b/benchmarks/test_scaling.py index 3b9a1e5..3245d1b 100644 --- a/benchmarks/test_scaling.py +++ b/benchmarks/test_scaling.py @@ -36,6 +36,7 @@ import collections import dataclasses import sys +import sysconfig import pytest from django.core.management import call_command @@ -46,7 +47,6 @@ ) 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.""" @@ -158,8 +158,10 @@ def test_drain_cpu_workload__benchmark(self, benchmark, parallelism, empty_queue "threads": parallelism.threads, "tasks": TASK_COUNT, "python": sys.version.split()[0], - "free_threaded_build": is_free_threaded_build(), - "gil_enabled": is_gil_enabled(), + "free_threaded_build": bool( + sysconfig.get_config_var("Py_GIL_DISABLED") + ), + "gil_enabled": getattr(sys, "_is_gil_enabled", lambda: True)(), } ) diff --git a/tests/test_executor.py b/tests/test_executor.py index ebdaf7f..a89a7a0 100644 --- a/tests/test_executor.py +++ b/tests/test_executor.py @@ -46,8 +46,6 @@ WorkerThread, configure_logging, handler, - is_free_threaded_build, - is_gil_enabled, warn_when_free_threading_is_unavailable, ) @@ -304,62 +302,53 @@ 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) + monkeypatch.setattr(sysconfig, "get_config_var", lambda name: 1) + monkeypatch.setattr(sys, "_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 + assert "may re-enable the GIL" 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) + monkeypatch.setattr(sysconfig, "get_config_var", lambda name: 1) + monkeypatch.setattr(sys, "_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 + assert "may re-enable the GIL" 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) + """Stay silent on an interpreter without the build flag.""" + monkeypatch.setattr(sysconfig, "get_config_var", lambda name: None) + monkeypatch.setattr(sys, "_is_gil_enabled", lambda: True) + + with caplog.at_level(logging.WARNING, logger="multiprocessing"): + warn_when_free_threading_is_unavailable() + + assert "may re-enable the GIL" not in caplog.text + + def test_warn_when_free_threading_is_unavailable__assume_gil_without_api( + self, monkeypatch, caplog + ): + """Assume the GIL is enabled when the private API is absent.""" + monkeypatch.setattr(sysconfig, "get_config_var", lambda name: 1) + monkeypatch.delattr(sys, "_is_gil_enabled", raising=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 + assert "may re-enable the GIL" in caplog.text def _wait_for_result(task_id: str) -> TaskResult: diff --git a/threadmill/executor.py b/threadmill/executor.py index 5edebfc..0228e49 100644 --- a/threadmill/executor.py +++ b/threadmill/executor.py @@ -100,16 +100,6 @@ def configure_logging(formatter: logging.Formatter) -> None: root_logger.addHandler(handler) -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. @@ -117,11 +107,12 @@ def warn_when_free_threading_is_unavailable() -> None: for the whole process, which silently costs the parallelism the build exists to provide. """ - if is_free_threaded_build() and is_gil_enabled(): + free_threaded_build = bool(sysconfig.get_config_var("Py_GIL_DISABLED")) + gil_enabled = getattr(sys, "_is_gil_enabled", lambda: True)() + if free_threaded_build and 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" + "Free-threaded interpreter is running with the GIL enabled. C extensions " + "that do not declare free-threading support may re-enable the GIL" ) @@ -179,8 +170,6 @@ 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 8d98217f5fba84502f68cb9d7819c0ffa6799e8f Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Thu, 8 Oct 2026 17:37:58 +0200 Subject: [PATCH 6/7] Mark the idle polling state private IdleBackoff sat next to threadmill.retry.ExponentialBackoff, which retries failed tasks rather than pacing idle polls. Nothing is expected to import it, so name it _IdleBackoff and say in its docstring what sets it apart. --- tests/backends/test_redis.py | 4 ++-- threadmill/backends/redis.py | 10 ++++++---- 2 files changed, 8 insertions(+), 6 deletions(-) diff --git a/tests/backends/test_redis.py b/tests/backends/test_redis.py index 9d623e9..8b60aac 100644 --- a/tests/backends/test_redis.py +++ b/tests/backends/test_redis.py @@ -34,9 +34,9 @@ TelemetryEvent, ) from threadmill.backends.redis import ( # noqa: E402 - IdleBackoff, RedisBroker, RedisTaskBackend, + _IdleBackoff, ) TELEMETRY_INTERVAL = datetime.timedelta(seconds=60) @@ -1834,7 +1834,7 @@ def test_idle_backoff__randomize_rotation_offset(self): 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]] = {} + observed: dict[str, list[_IdleBackoff]] = {} def observe(name: str) -> None: first = backend._idle_backoff diff --git a/threadmill/backends/redis.py b/threadmill/backends/redis.py index ae0ada8..5eb7d1d 100644 --- a/threadmill/backends/redis.py +++ b/threadmill/backends/redis.py @@ -51,11 +51,13 @@ def _parse_lease_started_at(value: str | None) -> datetime.datetime | None: @dataclasses.dataclass(kw_only=True, slots=True) -class IdleBackoff: +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. + misses never reset or double another thread's delay. Private, so it is not + mistaken for a public backoff like ``threadmill.retry.ExponentialBackoff``, + which retries failed tasks rather than pacing idle polls. """ miss_count: int = 0 @@ -245,11 +247,11 @@ def __init__(self, alias: str, params: dict) -> None: self._acknowledge_script = self.client.register_script(self.ACKNOWLEDGE_SCRIPT) @property - def _idle_backoff(self) -> IdleBackoff: + 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( + backoff = _IdleBackoff( rotation_offset=random.randrange(len(self.queues)) # noqa: S311 ) self._idle_backoffs.backoff = backoff From 5ea303f53d8ea7a268286e78a514c9a87873a2c2 Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Thu, 8 Oct 2026 17:38:01 +0200 Subject: [PATCH 7/7] Benchmark every queue with the same worker pool The chart compared four threads on a free-threaded build against one thread everywhere else, and needed a footnote to explain the difference away. - Run four threads wherever the worker supports threads: celery on its thread pool, dramatiq and huey by flag, threadmill by threads. Every prefetch window stays at READ_AHEAD messages. - Keep one thread for django-tasks-db and django-tasks-rq, which ship single-threaded workers and cannot be told otherwise. - Take the free-threading row from a 3.14t run and leave the rest of the chart on the GIL build, so both Threadmill bars are the same pool on two interpreters. chart.py accepts the second run as another argument. - Drop the footnote, along with the theme colour only it used. - Deepen the threadmill queues to 120,000 tasks. At 60,000 the marginal drain of the free-threading row is one to two seconds, so the one-second quantization of the worker start and stop decided the reported rate. - Re-measure and update the README alt text. --- README.md | 2 +- benchmarks/chart.py | 120 ++++++++++++----------- benchmarks/test_backends.py | 110 ++++++++++++--------- docs/images/backend-comparison-dark.svg | 38 ++++--- docs/images/backend-comparison-light.svg | 38 ++++--- 5 files changed, 166 insertions(+), 142 deletions(-) diff --git a/README.md b/README.md index 4b17ff6..f4630a6 100644 --- a/README.md +++ b/README.md @@ -19,7 +19,7 @@ - Tasks per second with one worker process each: Threadmill (free threading) 60,259, Threadmill 12,021, dramatiq 6,967, huey 5,389, celery 2,186, django-tasks-db 2,056, django-tasks-rq 75. + 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.

diff --git a/benchmarks/chart.py b/benchmarks/chart.py index aed3c04..d08a196 100644 --- a/benchmarks/chart.py +++ b/benchmarks/chart.py @@ -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 @@ -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' @@ -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.""" @@ -62,7 +63,6 @@ class Theme: border="#e5e7eb", ink="#111827", muted="#6b7280", - faint="#9ca3af", accent="#4f46e5", accent_bar="#4f46e5", ) @@ -72,7 +72,6 @@ class Theme: border="#30363d", ink="#e6edf3", muted="#8b949e", - faint="#6e7681", accent="#818cf8", accent_bar="#6366f1", ) @@ -88,8 +87,17 @@ class Theme: 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.""" +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) @@ -136,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 ( @@ -148,35 +184,18 @@ 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 process each: " + ", ".join( - f"{result.name} {result.throughput:,.0f}" for result in results + f"{display_name(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. - - Two lines because the whole footnote no longer fits the canvas at a legible - size. 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, -rq and huey read one task 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] @@ -191,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' str: text( LABEL_X, center + 5, - result.name, + display_name(result.name), theme=theme, size=14, weight=700 if is_fastest else 400, @@ -245,35 +263,25 @@ def build_chart(results: list[QueueResult], theme: Theme) -> str: ) ) - 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 calibrated from the 11.5px hand check below (right edge - # 863 of 900 for 141 characters, about 5.9px per character); the - # longest line now reaches about 626 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) if __name__ == "__main__": - if len(sys.argv) != 2: - raise SystemExit(f"usage: {sys.argv[0]} ") + if len(sys.argv) not in {2, 3}: + raise SystemExit( + f"usage: {sys.argv[0]} []" + ) 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)" ) diff --git a/benchmarks/test_backends.py b/benchmarks/test_backends.py index 0101e06..cbf59ae 100644 --- a/benchmarks/test_backends.py +++ b/benchmarks/test_backends.py @@ -18,7 +18,10 @@ the task count of the queue by the difference. The result is the marginal throughput of a busy queue. -Every queue runs one worker process and one thread. Each queue reads ``READ_AHEAD`` +Every queue runs one worker process and ``WORKER_THREAD_COUNT`` threads where its +worker supports threads, so the bars rank the queues and their interpreters rather +than the size of their pools. django-tasks-db and django-tasks-rq ship +single-threaded workers and run one thread each. Each queue reads ``READ_AHEAD`` messages ahead where the queue has such a setting. The numbers therefore rank the queues and not their polling strategies. The consumer of Celery blocks on an empty queue, so its window costs no sleep for each message. Dramatiq polls instead and @@ -40,9 +43,10 @@ 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. +more with the same ``WORKER_THREAD_COUNT`` threads as every other queue. That row is +listed only there, because on any other interpreter its threads run one at a time +and the drain would only repeat the single-threaded rate. The two Threadmill rows +are therefore the same worker configuration on two interpreters. 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 @@ -57,6 +61,7 @@ import os import subprocess import sys +import sysconfig import tempfile import time import typing @@ -82,7 +87,6 @@ from benchmarks.dramatiq_app import dramatiq_echo, dramatiq_mark_processed from benchmarks.huey_app import huey_echo, huey_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.""" @@ -108,13 +112,13 @@ and has no read-ahead setting either. """ -FREE_THREADING_THREAD_COUNT = 4 -"""Threads the free-threading threadmill worker runs, where the build allows it. +WORKER_THREAD_COUNT = 4 +"""Threads every worker runs, where its queue supports threads. -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. +The same pool size for every queue, so the comparison isolates the queue and its +interpreter rather than the pool. Deliberately a fixed small number rather than +every core, so the bars state what a given pool does and not what this machine +happens to have. """ CELERY_WORKER = ( @@ -124,21 +128,22 @@ "-A", "benchmarks.celery_app:celery_app", "worker", - "--pool=solo", - f"--prefetch-multiplier={READ_AHEAD}", + "--pool=threads", + f"--concurrency={WORKER_THREAD_COUNT}", + f"--prefetch-multiplier={READ_AHEAD // WORKER_THREAD_COUNT}", "--loglevel=WARNING", "--without-gossip", "--without-mingle", "--without-heartbeat", ) -"""Celery worker running as one process with one thread, ``READ_AHEAD`` messages ahead. +"""Celery worker running as one process with ``WORKER_THREAD_COUNT`` threads, ``READ_AHEAD`` messages ahead. The default prefork pool crashes on CPython 3.14. The pool child loses the task -handler state that it expects. The worker therefore runs on the solo pool. With -one concurrent task, ``--prefetch-multiplier={READ_AHEAD}`` sets the prefetch -count to ``READ_AHEAD``, which is the benchmark rate. The Redis consumer of -Celery blocks while its queue is empty, so the prefetch adds no sleep for each -message. +handler state that it expects. The worker therefore runs on the thread pool. Celery +multiplies the concurrency by ``--prefetch-multiplier``, so the multiplier +``READ_AHEAD // WORKER_THREAD_COUNT`` holds the window at ``READ_AHEAD`` messages. +The Redis consumer of Celery blocks while its queue is empty, so the prefetch adds +no sleep for each message. """ DRAMATIQ_WORKER = ( @@ -149,15 +154,16 @@ "--processes", "1", "--threads", - "1", + f"{WORKER_THREAD_COUNT}", ) -"""dramatiq worker running as one process with one thread, ``READ_AHEAD`` messages ahead. +"""dramatiq worker running as one process with ``WORKER_THREAD_COUNT`` threads, ``READ_AHEAD`` messages ahead. The Redis broker polls and does not block. Its consumer fetches only while fewer messages than its read-ahead are unacked. When that window is full, the consumer sleeps a jittered backoff of 5 to 10 ms and then polls again. The command line has no read-ahead flag, so the worker environment carries -``dramatiq_queue_prefetch``. The benchmark sets this variable to ``READ_AHEAD``. +``dramatiq_queue_prefetch``. The benchmark sets this variable to ``READ_AHEAD``, +which is an absolute window rather than one for each thread. """ HUEY_WORKER = ( @@ -165,13 +171,13 @@ "-m", "huey.bin.huey_consumer", "benchmarks.huey_app.huey_app", - "--workers=1", + f"--workers={WORKER_THREAD_COUNT}", "--worker-type=thread", "--no-periodic", "--quiet", "--graceful-signal=TERM", ) -"""huey consumer running as one process with one thread, one message at a time. +"""huey consumer running as one process with ``WORKER_THREAD_COUNT`` threads, one message at a time. The Redis storage of huey blocks on an empty queue and pops a single message. It has no read-ahead setting. ``--quiet`` matches the log level of the other @@ -216,13 +222,14 @@ class QueueUnderTest: def drain_with_threadmill_worker() -> None: - """Process every queued task with a single threadmill worker process.""" + """Process every queued task with one threadmill worker process on ``WORKER_THREAD_COUNT`` threads.""" call_command( "threadmill", "worker", backend=DEFAULT_TASK_BACKEND_ALIAS, queues=[DEFAULT_TASK_QUEUE_NAME], workers=1, + threads=WORKER_THREAD_COUNT, prefetch_count=READ_AHEAD, exit_empty=True, verbosity=0, @@ -230,7 +237,7 @@ def drain_with_threadmill_worker() -> None: def drain_with_threadmill_worker_no_prefetch() -> None: - """Process every queued task with one threadmill worker reading one at a time. + """Process every queued task with the same pool reading one task at a time. The no-prefetch ablation of the entry above, so the read-ahead cost can be subtracted. @@ -241,6 +248,7 @@ def drain_with_threadmill_worker_no_prefetch() -> None: backend=DEFAULT_TASK_BACKEND_ALIAS, queues=[DEFAULT_TASK_QUEUE_NAME], workers=1, + threads=WORKER_THREAD_COUNT, prefetch_count=1, exit_empty=True, verbosity=0, @@ -248,12 +256,13 @@ 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. + """Process every queued task with the same pool on a free-threaded interpreter. - 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. + Same one process and ``WORKER_THREAD_COUNT`` threads as every other queue, so + the two Threadmill rows differ in the interpreter alone. Only meaningful on a + free-threaded interpreter, where the threads run at the same time. The queue is + only listed under test there, so this drain never reports the single-threaded + rate of a GIL build as if it were a parallel one. """ call_command( "threadmill", @@ -261,7 +270,7 @@ def drain_with_threadmill_free_threading_worker() -> None: backend=DEFAULT_TASK_BACKEND_ALIAS, queues=[DEFAULT_TASK_QUEUE_NAME], workers=1, - threads=FREE_THREADING_THREAD_COUNT, + threads=WORKER_THREAD_COUNT, prefetch_count=READ_AHEAD, exit_empty=True, verbosity=0, @@ -296,13 +305,13 @@ def drain_with_django_tasks_rq_worker() -> None: def drain_with_celery_worker() -> None: - """Process every queued task with a single-process Celery worker.""" + """Process every queued task with a Celery worker on ``WORKER_THREAD_COUNT`` threads.""" celery_mark_processed.delay() drain_with_subprocess_worker(CELERY_WORKER) def drain_with_dramatiq_worker() -> None: - """Process every queued task with a single-process, single-thread dramatiq worker.""" + """Process every queued task with a dramatiq worker on ``WORKER_THREAD_COUNT`` threads.""" dramatiq_mark_processed.send() drain_with_subprocess_worker( DRAMATIQ_WORKER, @@ -311,7 +320,7 @@ def drain_with_dramatiq_worker() -> None: def drain_with_huey_worker() -> None: - """Process every queued task with a single-thread huey consumer.""" + """Process every queued task with a huey consumer on ``WORKER_THREAD_COUNT`` threads.""" huey_mark_processed() drain_with_subprocess_worker(HUEY_WORKER) @@ -387,7 +396,10 @@ def django_task_enqueuer( def enqueue(count: int) -> TaskResult: """Accept `count` echo tasks, returning the newest task result.""" - return [task.enqueue(index) for index in range(count)][-1] + task_result = task.enqueue(0) + for index in range(1, count): + task_result = task.enqueue(index) + return task_result def verify_processed(enqueued_task_result: TaskResult | None) -> None: """Assert that the backend executed the benchmark tasks. @@ -436,15 +448,24 @@ def django_task_backend( ) -THREADMILL_TASK_COUNT = 60_000 -"""Tasks threadmill queues, so its marginal drain outruns the one-second fixed cost. +THREADMILL_TASK_COUNT = 120_000 +"""Tasks every threadmill queue holds, so its marginal drain outruns the one-second fixed cost. + +Threadmill drains this queue in about 14 seconds with four threads on a GIL build +and in about 4 seconds on a free-threaded one. The fixed cost of a cold worker start +and stop is quantized to about a second. A shallower queue therefore keeps the +prefetch comparison inside that step. -Threadmill drains a queue in about 3 seconds for each 20,000 tasks. The fixed cost -of a cold worker start and stop is quantized to about a second. A shallower queue -therefore keeps the prefetch comparison inside that step. +The same depth for both interpreters keeps the two Threadmill rows comparable. A +deeper queue measures a slower rate per task, because a larger keyspace costs the +broker more, so a depth that differs between the rows would show up as a difference +between the interpreters. """ -THREADS_RUN_IN_PARALLEL = is_free_threaded_build() and not is_gil_enabled() +THREADS_RUN_IN_PARALLEL = ( + bool(sysconfig.get_config_var("Py_GIL_DISABLED")) + and not getattr(sys, "_is_gil_enabled", lambda: True)() +) """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 @@ -466,9 +487,8 @@ def django_task_backend( ) """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. +On any other interpreter the drain behind this queue runs the same pool with one +thread at a time, so it would report a free-threading result that found no speedup. """ WORKER_QUEUES = ( diff --git a/docs/images/backend-comparison-dark.svg b/docs/images/backend-comparison-dark.svg index 7f3a279..706bf63 100644 --- a/docs/images/backend-comparison-dark.svg +++ b/docs/images/backend-comparison-dark.svg @@ -1,29 +1,27 @@ - + - + Queue throughput -5,000–60,000 trivial tasks per queue · one worker process each · higher is better -Threadmill (free threading) +5,000–120,000 trivial tasks per queue · one worker process each · higher is better +Threadmill -60,259/s -Threadmill - -12,021/s +19,954/s +Threadmill (GIL) + +9,215/s dramatiq - -6,967/s + +7,251/s huey - -5,389/s + +6,254/s celery - -2,186/s + +2,497/s django-tasks-db - -2,056/s + +1,814/s django-tasks-rq - -75/s -One process and one thread each, except Threadmill on a free-threaded build, which runs four. -Threadmill, celery and dramatiq read 128 ahead; django-tasks-db, -rq and huey read one task at a time. + +84/s diff --git a/docs/images/backend-comparison-light.svg b/docs/images/backend-comparison-light.svg index 36a2b0f..73ad649 100644 --- a/docs/images/backend-comparison-light.svg +++ b/docs/images/backend-comparison-light.svg @@ -1,29 +1,27 @@ - + - + Queue throughput -5,000–60,000 trivial tasks per queue · one worker process each · higher is better -Threadmill (free threading) +5,000–120,000 trivial tasks per queue · one worker process each · higher is better +Threadmill -60,259/s -Threadmill - -12,021/s +19,954/s +Threadmill (GIL) + +9,215/s dramatiq - -6,967/s + +7,251/s huey - -5,389/s + +6,254/s celery - -2,186/s + +2,497/s django-tasks-db - -2,056/s + +1,814/s django-tasks-rq - -75/s -One process and one thread each, except Threadmill on a free-threaded build, which runs four. -Threadmill, celery and dramatiq read 128 ahead; django-tasks-db, -rq and huey read one task at a time. + +84/s