diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 58c8d91..cda143e 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -43,6 +43,31 @@ jobs: - uses: codecov/codecov-action@v7 with: token: ${{ secrets.CODECOV_TOKEN }} + pytest-free-threaded: + name: Pytest free-threaded + strategy: + matrix: + os: + - "ubuntu-latest" + python-version: + - "3.14t" + django-version: + - "6.1" + runs-on: ${{ matrix.os }} + services: + redis: + image: redis + ports: + - 6379:6379 + options: --entrypoint redis-server + env: + REDIS_URL: redis:///0 + steps: + - uses: actions/checkout@v7 + - uses: astral-sh/setup-uv@v7 + with: + python-version: ${{ matrix.python-version }} + - run: uv run --with django~=${{ matrix.django-version }}.0 pytest -m "not benchmark" pytest-extras: name: Pytest strategy: 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 a689c88..f4630a6 100644 --- a/README.md +++ b/README.md @@ -19,7 +19,7 @@ - Tasks per second with one worker: threadmill 9,960, dramatiq 7,177, huey 6,038, celery 2,280, django-tasks-db 2,065, django-tasks-rq 82. + 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.

@@ -78,13 +78,24 @@ The workers are inspired by Gunicorn, and the CLI is very similar. #### Utilization -Depending on your workload, you can tweak the number of processes and threads. -Processes allow for parallel compute (no GIL) while threads are great for low-memory concurrent IO. +Depending on your workload, you can tweak the number of worker processes and the number of threads per process. Processes always run in parallel, because each one has its own interpreter. Threads share memory, so they are cheap, but whether they run in parallel depends on your interpreter: + +- On a regular GIL build, Python runs in one thread at a time. Threads still overlap waiting for IO, but CPU-bound tasks do not speed up. Raising the thread count far above the queue depth can also leave processes idle, because each process prefetches a batch sized by its thread count. +- On a [free-threaded build](https://peps.python.org/pep-0779/), threads run truly in parallel, so CPU-bound tasks scale across cores without the cost of extra processes. + +A pool on a free-threaded interpreter therefore reaches the same throughput with one process running many threads as with many processes running one thread each, while using a single Redis connection pool and one copy of your application state: ```console -uv run manage.py threadmill worker --workers 4 --threads 2 +uv run manage.py threadmill worker --workers 1 --threads 8 ``` +Threads all live in one process, so a crash or a `--max-tasks` recycle takes every thread down at once. One process per core, the default, confines that to a single process. Prefer threads when memory matters more than that isolation, and keep the default when it does not. + +A free-threaded interpreter can be slower than a regular one for single-threaded work, so a GIL build stays the better choice for IO-bound tasks. Pick the interpreter that fits your workload rather than assuming free threading is an upgrade. + +> [!WARNING] +> A C extension that does not declare free-threading support re-enables the GIL for the whole process, which silently costs you all thread parallelism. `hiredis` is a common offender: installing `redis[hiredis]` enables it on import. The worker logs a warning when it detects that the GIL was re-enabled. + #### Health If your tasks leak memory, you can recycle (restart) the workers after a certain number of tasks have been processed: 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/chart.py b/benchmarks/chart.py index f5ad3c2..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,12 +72,11 @@ class Theme: border="#30363d", ink="#e6edf3", muted="#8b949e", - faint="#6e7681", accent="#818cf8", accent_bar="#6366f1", ) -DIAGNOSTIC_QUEUES = frozenset({"threadmill (no prefetch)"}) +DIAGNOSTIC_QUEUES = frozenset({"Threadmill (no prefetch)"}) """Queues that the benchmark measures and the chart leaves out. The harness runs threadmill twice to compare its prefetch buffer with a single @@ -85,6 +84,21 @@ 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.""" + +DISPLAY_NAMES = { + "Threadmill (free threading)": "Threadmill", + "Threadmill": "Threadmill (GIL)", +} +"""Chart labels for the threadmill rows, keyed by benchmark queue name. + +The benchmark names identify a row across runs, so they stay as they are and the +chart translates them for the reader. The free-threading result is threadmill's +headline row and carries the plain name, and the row measured next to the other +queues is labelled as the GIL baseline. +""" + @dataclasses.dataclass(frozen=True, kw_only=True, slots=True) class QueueResult: @@ -130,9 +144,37 @@ def read_results(json_path: pathlib.Path) -> list[QueueResult]: task_count=process["extra_info"]["tasks"], ) ) + return rank(results) + + +def rank(results: list[QueueResult]) -> list[QueueResult]: + """Return the results, fastest first.""" return sorted(results, key=lambda result: result.throughput, reverse=True) +def with_free_threading_result( + results: list[QueueResult], free_threading_results: list[QueueResult] +) -> list[QueueResult]: + """Return the results with the free-threading row taken from another run. + + The free-threading row is measured on a free-threaded interpreter, while every + other row is measured on the GIL build the chart ranks, so the two runs are + merged. The remaining rows of the second run repeat queues that the first run + already measured and are ignored. + """ + free_threading = [ + result + for result in free_threading_results + if result.name == FREE_THREADING_QUEUE + ] + if not free_threading: + raise SystemExit(f"no {FREE_THREADING_QUEUE} row in the free-threaded run") + return rank( + [result for result in results if result.name != FREE_THREADING_QUEUE] + + free_threading + ) + + def text(x, y, content, *, theme, size=13, fill=None, weight=400, anchor="start"): """Render one SVG text element.""" return ( @@ -142,10 +184,15 @@ def text(x, y, content, *, theme, size=13, fill=None, weight=400, anchor="start" ) +def display_name(queue_name: str) -> str: + """Return the chart label of a queue.""" + return DISPLAY_NAMES.get(queue_name, queue_name) + + def describe(results: list[QueueResult]) -> str: """Return a sentence describing the throughput of every queue.""" - return "Tasks per second with one worker: " + ", ".join( - f"{result.name} {result.throughput:,.0f}" for result in results + return "Tasks per second with one worker process each: " + ", ".join( + f"{display_name(result.name)} {result.throughput:,.0f}" for result in results ) @@ -163,8 +210,7 @@ def build_chart(results: list[QueueResult], theme: Theme) -> str: row_centers = [ FIRST_ROW_CENTER + index * ROW_HEIGHT for index in range(len(results)) ] - footnote_y = row_centers[-1] + FOOTNOTE_GAP - height = footnote_y + BOTTOM_PADDING + height = row_centers[-1] + BOTTOM_PADDING parts = [ f' 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, @@ -192,7 +238,7 @@ def build_chart(results: list[QueueResult], theme: Theme) -> str: text( LABEL_X, center + 5, - result.name, + display_name(result.name), theme=theme, size=14, weight=700 if is_fastest else 400, @@ -217,32 +263,25 @@ def build_chart(results: list[QueueResult], theme: Theme) -> str: ) ) - parts.append( - text( - 28, - footnote_y, - # joe: width checked by hand (right edge 863 of 900 at 11.5px); add a - # width guard if the canvas width or the font stack changes. - "One process and one thread each. Threadmill, celery and dramatiq read " - "128 ahead; django-tasks-db, -rq and huey read one task at a time.", - theme=theme, - size=11.5, - fill=theme.faint, - ) - ) parts.append("") 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/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_backends.py b/benchmarks/test_backends.py index b8cf81c..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 @@ -39,6 +42,12 @@ 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 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 sentinel task last and waits for it. That wait proves that the queue was drained. @@ -52,6 +61,7 @@ import os import subprocess import sys +import sysconfig import tempfile import time import typing @@ -102,6 +112,15 @@ and has no read-ahead setting either. """ +WORKER_THREAD_COUNT = 4 +"""Threads every worker runs, where its queue supports threads. + +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 = ( sys.executable, "-m", @@ -109,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 = ( @@ -134,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 = ( @@ -150,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 @@ -201,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, @@ -215,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. @@ -226,12 +248,35 @@ 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, ) +def drain_with_threadmill_free_threading_worker() -> None: + """Process every queued task with the same pool on a free-threaded interpreter. + + 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", + "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, + ) + + def drain_with_django_tasks_db_worker() -> None: """Process every queued task with the django-tasks-db worker.""" call_command( @@ -260,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, @@ -275,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) @@ -351,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. @@ -400,23 +448,59 @@ 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 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. +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. + +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 = ( + 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 +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 runs the same pool with one +thread at a time, so it would report 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/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() diff --git a/benchmarks/test_scaling.py b/benchmarks/test_scaling.py new file mode 100644 index 0000000..3245d1b --- /dev/null +++ b/benchmarks/test_scaling.py @@ -0,0 +1,176 @@ +"""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. + +Measured on 16 CPU-bound tasks of about a second each, worker start cost included. +On a free-threaded build one process with four threads reaches the same +throughput as four processes with one thread each, about 2.7x, and four +processes with four threads match them: worker count sets the parallelism, and +threads reach it with less memory. On a GIL build extra threads never help, and +they cost once a pool holds more worker threads in total than the queue holds +tasks. A process requests a prefetch batch sized by its thread count, so one +process can take the whole queue and leave its neighbours idle. + +That is why the 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 collections +import dataclasses +import sys +import sysconfig + +import pytest +from django.core.management import call_command +from django.tasks import ( + DEFAULT_TASK_BACKEND_ALIAS, + TaskResultStatus, + task_backends, +) + +from tests.testapp.tasks import compute_workload + +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.""" + +SCALING_QUEUE_NAME = "scaling" +"""Queue dedicated to this benchmark. + +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. +""" + +scaling_workload = dataclasses.replace(compute_workload, queue_name=SCALING_QUEUE_NAME) +"""The CPU-bound workload task, pinned to the benchmark's own queue.""" + + +@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(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 + enqueued_ids: list[str] = dataclasses.field(default_factory=list, init=False) + + def drain(self) -> None: + """Queue the workload, then process it with the configured worker pool.""" + queued_task = scaling_workload.using(backend=DEFAULT_TASK_BACKEND_ALIAS) + self.enqueued_ids.extend( + queued_task.enqueue().id for _ in range(self.task_count) + ) + call_command( + "threadmill", + "worker", + backend=DEFAULT_TASK_BACKEND_ALIAS, + queues=[SCALING_QUEUE_NAME], + workers=self.parallelism.workers, + threads=self.parallelism.threads, + exit_empty=True, + verbosity=0, + ) + + def verify(self) -> None: + """Assert every queued task succeeded, so a fast drain cannot be a lost task.""" + backend = task_backends[DEFAULT_TASK_BACKEND_ALIAS] + statuses = collections.Counter( + backend.get_result(task_id).status for task_id in self.enqueued_ids + ) + assert statuses[TaskResultStatus.SUCCESSFUL] == len(self.enqueued_ids), ( + f"drained {statuses[TaskResultStatus.SUCCESSFUL]} of" + f" {len(self.enqueued_ids)} tasks: {statuses}" + ) + + +@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.""" + workload = CpuWorkload(parallelism=parallelism) + benchmark.extra_info.update( + { + "workers": parallelism.workers, + "threads": parallelism.threads, + "tasks": TASK_COUNT, + "python": sys.version.split()[0], + "free_threaded_build": bool( + sysconfig.get_config_var("Py_GIL_DISABLED") + ), + "gil_enabled": getattr(sys, "_is_gil_enabled", lambda: True)(), + } + ) + + benchmark.pedantic( + workload.drain, + rounds=MEASUREMENT_ROUNDS, + iterations=1, + warmup_rounds=0, + ) + + # Outside the timed region: every round must have run every task. + workload.verify() diff --git a/docs/images/backend-comparison-dark.svg b/docs/images/backend-comparison-dark.svg index 41f1146..706bf63 100644 --- a/docs/images/backend-comparison-dark.svg +++ b/docs/images/backend-comparison-dark.svg @@ -1,25 +1,27 @@ - + - + Queue throughput -5,000–60,000 trivial tasks per queue · one worker process, one thread · higher is better -threadmill +5,000–120,000 trivial tasks per queue · one worker process each · higher is better +Threadmill -9,960/s -dramatiq - -7,177/s -huey - -6,038/s -celery - -2,280/s -django-tasks-db - -2,065/s -django-tasks-rq - -82/s -One process and one thread each. Threadmill, celery and dramatiq read 128 ahead; django-tasks-db, -rq and huey read one task at a time. +19,954/s +Threadmill (GIL) + +9,215/s +dramatiq + +7,251/s +huey + +6,254/s +celery + +2,497/s +django-tasks-db + +1,814/s +django-tasks-rq + +84/s diff --git a/docs/images/backend-comparison-light.svg b/docs/images/backend-comparison-light.svg index 1dac111..73ad649 100644 --- a/docs/images/backend-comparison-light.svg +++ b/docs/images/backend-comparison-light.svg @@ -1,25 +1,27 @@ - + - + Queue throughput -5,000–60,000 trivial tasks per queue · one worker process, one thread · higher is better -threadmill +5,000–120,000 trivial tasks per queue · one worker process each · higher is better +Threadmill -9,960/s -dramatiq - -7,177/s -huey - -6,038/s -celery - -2,280/s -django-tasks-db - -2,065/s -django-tasks-rq - -82/s -One process and one thread each. Threadmill, celery and dramatiq read 128 ahead; django-tasks-db, -rq and huey read one task at a time. +19,954/s +Threadmill (GIL) + +9,215/s +dramatiq + +7,251/s +huey + +6,254/s +celery + +2,497/s +django-tasks-db + +1,814/s +django-tasks-rq + +84/s diff --git a/pyproject.toml b/pyproject.toml index dbc1402..9289dee 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 cbd5da1..8b60aac 100644 --- a/tests/backends/test_redis.py +++ b/tests/backends/test_redis.py @@ -5,6 +5,7 @@ import json import logging import queue +import threading import time import typing from dataclasses import replace @@ -35,6 +36,7 @@ from threadmill.backends.redis import ( # noqa: E402 RedisBroker, RedisTaskBackend, + _IdleBackoff, ) TELEMETRY_INTERVAL = datetime.timedelta(seconds=60) @@ -1637,7 +1639,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]) @@ -1666,7 +1668,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() @@ -1679,7 +1681,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: @@ -1691,7 +1693,7 @@ def test_acquire__wraps_rotation_offset_through_every_queue(self): *queue_names, timeout=datetime.timedelta(seconds=1) ) acquired.append(task_result) - 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", @@ -1715,7 +1717,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( @@ -1726,7 +1728,7 @@ def test_acquire__keep_rotation_on_idle_polls(self): assert len(recorder.sent_args) > 1 assert {sent_args[-2] for sent_args in recorder.sent_args} == {"5"} assert {sent_args[-1] for sent_args in recorder.sent_args} == {"1"} - assert backend._rotation_offset == 5 + assert backend._idle_backoff.rotation_offset == 5 finally: backend.close() @@ -1793,7 +1795,7 @@ def test_acquire__spreads_batch_across_queues(self): "acquire_batch_round_robin_test", queues=["compute", "io"], ) - backend._rotation_offset = 0 + backend._idle_backoff.rotation_offset = 0 try: for value in range(2): backend.enqueue(replace(echo, queue_name="compute"), args=[value]) @@ -1810,11 +1812,11 @@ def test_acquire__spreads_batch_across_queues(self): "compute", "io", ] - assert backend._rotation_offset == 1 + assert backend._idle_backoff.rotation_offset == 1 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 = [ @@ -1822,7 +1824,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 45bccfb..a89a7a0 100644 --- a/tests/test_executor.py +++ b/tests/test_executor.py @@ -6,6 +6,7 @@ import multiprocessing import queue import sys +import sysconfig import threading import time import uuid @@ -19,6 +20,7 @@ default_task_backend, task, ) +from django.tasks.exceptions import TaskResultDoesNotExist from django.utils import timezone from tests.testapp.tasks import ( @@ -29,6 +31,7 @@ count_users, echo, log_message, + record_execution, ) from threadmill.backends.base import ( Broker, @@ -43,6 +46,7 @@ WorkerThread, configure_logging, handler, + warn_when_free_threading_is_unavailable, ) @@ -295,6 +299,74 @@ def test_configure_logging__keeps_placeholder_loggers(self): ) +class TestFreeThreadingDiagnostic: + """Tests for the interpreter free-threading diagnostic.""" + + 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(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 "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(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 "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 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 "may re-enable the GIL" in caplog.text + + +def _wait_for_result(task_id: str) -> TaskResult: + """Poll for a task result until the worker persists it. + + A spawned worker boots Django before it can process anything, so a fixed + sleep races the worker instead of waiting for it. + """ + deadline = time.monotonic() + 30 + while True: + try: + return default_task_backend.get_result(task_id) + except TaskResultDoesNotExist: + if time.monotonic() >= deadline: + raise + time.sleep(0.01) + + class TestTaskExecutor: """Tests for the TaskExecutor dataclass and its methods.""" @@ -425,6 +497,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"]) @@ -476,11 +574,10 @@ def test_run__executes_model_task_in_spawned_worker(self): ) run_thread = threading.Thread(target=executor.run, daemon=True) run_thread.start() - time.sleep(3) + result = _wait_for_result(enqueued.id) executor.shutdown() run_thread.join(timeout=5) assert not run_thread.is_alive() - result = default_task_backend.get_result(enqueued.id) assert result.status == TaskResultStatus.SUCCESSFUL finally: multiprocessing.set_start_method(original_start_method, force=True) diff --git a/tests/testapp/settings.py b/tests/testapp/settings.py index db8e734..aab60e6 100644 --- a/tests/testapp/settings.py +++ b/tests/testapp/settings.py @@ -94,7 +94,7 @@ TASKS = { DEFAULT_TASK_BACKEND_ALIAS: { "BACKEND": "threadmill.backends.redis.RedisTaskBackend", - "QUEUES": [DEFAULT_TASK_QUEUE_NAME, "compute", "io", "memory"], + "QUEUES": [DEFAULT_TASK_QUEUE_NAME, "compute", "io", "memory", "scaling"], "REDIS_URL": REDIS_URL, "OPTIONS": { "max_connections": 10, diff --git a/tests/testapp/tasks.py b/tests/testapp/tasks.py index cd6fa2a..d8ed7a2 100644 --- a/tests/testapp/tasks.py +++ b/tests/testapp/tasks.py @@ -31,6 +31,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 5737e5c..5eb7d1d 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 @@ -49,6 +50,20 @@ def _parse_lease_started_at(value: str | None) -> datetime.datetime | None: return None +@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. 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 + rotation_offset: int = 0 + + class RedisBroker(Broker): """Background maintenance broker for the Redis backend.""" @@ -227,11 +242,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.""" @@ -330,6 +355,7 @@ def acquire( ) ] + idle_backoff = self._idle_backoff while True: now = timezone.now() now_ms = now.timestamp() * 1000 @@ -344,12 +370,14 @@ def acquire( str(len(queue_names)), worker, str(int(self.lease_ttl.total_seconds() * 1000)), - str(self._rotation_offset), + str(idle_backoff.rotation_offset), str(count), ], ): - 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._apply_lease( self.deserialize_task_result(data), @@ -370,11 +398,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) @staticmethod diff --git a/threadmill/executor.py b/threadmill/executor.py index 71615f0..0228e49 100644 --- a/threadmill/executor.py +++ b/threadmill/executor.py @@ -10,6 +10,7 @@ import random import socket import sys +import sysconfig import threading import time import typing @@ -99,6 +100,22 @@ def configure_logging(formatter: logging.Formatter) -> None: root_logger.addHandler(handler) +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. + """ + 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. C extensions " + "that do not declare free-threading support may re-enable the GIL" + ) + + @dataclasses.dataclass(kw_only=True, slots=True) class TaskExecutor: """Tasks consumed from shared joinable queues via process and thread pools.""" @@ -153,6 +170,7 @@ def create_worker_process(self) -> WorkerProcess: def run(self) -> None: """Start consuming tasks until shutdown is requested.""" configure_logging(self.log_formatter) + warn_when_free_threading_is_unavailable() self.worker_processes = [ self.create_worker_process() for _ in range(self.process_count) ]