From 8db8be8cff260eaf9bc141d20a1a33c3b5f12072 Mon Sep 17 00:00:00 2001 From: Johannes Maron Date: Wed, 7 Oct 2026 16:08:21 +0200 Subject: [PATCH] Measure django-tasks-rq and dramatiq in the queue benchmark django-tasks-redis leaves the comparison for django-tasks-rq, the RQ backend of the django-tasks family, and dramatiq joins it. Both are measured on the same trivial echo task and the same queue as the rest: - django-tasks-rq drains with `rqworker --burst` and its own Job class. RQ's worker forks a work horse per job, so its 80 tasks/s includes that fork, and its enqueue path costs 793 tasks/s against threadmill's 7,969. - dramatiq and celery read four messages per worker thread, celery's shipped prefetch multiplier. Dramatiq's Redis consumer polls rather than blocks and sleeps a jittered backoff once its window fills, so 669 tasks/s still measures that backoff more than the queue. The same worker drains 451 tasks/s at two messages, 1,983 at 16 and 5,658 at 64, where the backoff stops setting the rate; its own default is two per worker thread. - threadmill, django-tasks-db and django-tasks-rq read one task at a time, because their workers block on an empty queue and gain nothing from a window. The status check compares by value, since django-tasks-rq returns the django-tasks backport's TaskResultStatus, a different enum with the same value. The chart grows with its row count, and the README alt text carries the new numbers. --- README.md | 2 +- benchmarks/chart.py | 23 +++-- benchmarks/dramatiq_app.py | 49 ++++++++++ benchmarks/test_backends.py | 115 +++++++++++++++++++---- docs/images/backend-comparison-dark.svg | 25 ++--- docs/images/backend-comparison-light.svg | 25 ++--- pyproject.toml | 3 +- tests/testapp/settings.py | 13 ++- 8 files changed, 202 insertions(+), 53 deletions(-) create mode 100644 benchmarks/dramatiq_app.py diff --git a/README.md b/README.md index bade689..9112ff7 100644 --- a/README.md +++ b/README.md @@ -19,7 +19,7 @@ - Tasks per second with one worker: threadmill 5,069, celery 2,128, django-tasks-db 1,997, django-tasks-redis 1,313. + Tasks per second with one worker: threadmill 5,023, celery 1,951, django-tasks-db 1,942, dramatiq 669, django-tasks-rq 80.

diff --git a/benchmarks/chart.py b/benchmarks/chart.py index 2c12c66..a19346c 100644 --- a/benchmarks/chart.py +++ b/benchmarks/chart.py @@ -16,12 +16,16 @@ DARK_THEME_PATH = IMAGE_DIRECTORY / "backend-comparison-dark.svg" WIDTH = 900 -HEIGHT = 332 LABEL_X = 180 PLOT_X0 = 200 PLOT_WIDTH = 560 +FIRST_ROW_CENTER = 110 +ROW_HEIGHT = 40 +FOOTNOTE_GAP = 38 +BOTTOM_PADDING = 64 + FONT = ( 'system-ui, -apple-system, "Segoe UI", Roboto, "Helvetica Neue", Arial, sans-serif' ) @@ -139,14 +143,18 @@ def build_chart(results: list[QueueResult], theme: Theme) -> str: """Return the chart as an SVG document drawn in the given theme.""" fastest = results[0] scale = PLOT_WIDTH / fastest.throughput - row_centers = [110 + index * 40 for index in range(len(results))] + 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 parts = [ - f'', "", - f'', text(28, 46, "Queue throughput", theme=theme, size=19, weight=700), text( @@ -195,8 +203,9 @@ def build_chart(results: list[QueueResult], theme: Theme) -> str: parts.append( text( 28, - row_centers[-1] + 38, - "One message in flight per worker — no queue reads ahead.", + footnote_y, + "threadmill, django-tasks-db and django-tasks-rq read one message at a " + "time; celery and dramatiq four.", theme=theme, size=11.5, fill=theme.faint, diff --git a/benchmarks/dramatiq_app.py b/benchmarks/dramatiq_app.py new file mode 100644 index 0000000..b3f474a --- /dev/null +++ b/benchmarks/dramatiq_app.py @@ -0,0 +1,49 @@ +"""The dramatiq broker the comparison benchmarks hand tasks to. + +The dramatiq worker CLI imports this module without setting up Django, so it +must not import Django or any Django application. The broker carries the Results +middleware so the echo actor stores its return value in Redis, the way the Celery +app's result backend does. The result backend names its keys +``dramatiq:results:::`` rather than the default bare MD5 +hash, so the benchmark cleanup's ``dramatiq:*`` pattern deletes them. +""" + +import os + +import dramatiq +import redis +from dramatiq.brokers.redis import RedisBroker +from dramatiq.results import Results +from dramatiq.results.backends.redis import RedisBackend + +REDIS_URL = os.environ.get("REDIS_URL", "redis://localhost:6379/0") + +PROCESSED_KEY = "benchmark:processed" +"""Key the sentinel task increments once every earlier task was processed.""" + +client = redis.Redis.from_url(REDIS_URL) + +redis_broker = RedisBroker(url=REDIS_URL) +# A greppable namespace the benchmark cleanup's "dramatiq:*" pattern matches; +# the backend's default key is a bare MD5 hash no pattern can name. +redis_broker.add_middleware( + Results( + backend=RedisBackend( + client=client, + namespace="dramatiq:results", + use_namespace_prefix_keys=True, + ), + ), +) + + +@dramatiq.actor(broker=redis_broker, store_results=True) +def dramatiq_echo(value): + """Return the given value.""" + return value + + +@dramatiq.actor(broker=redis_broker) +def dramatiq_mark_processed(): + """Record that every earlier task in the queue has been processed.""" + client.incr(PROCESSED_KEY) diff --git a/benchmarks/test_backends.py b/benchmarks/test_backends.py index 4313a71..0ff4b95 100644 --- a/benchmarks/test_backends.py +++ b/benchmarks/test_backends.py @@ -12,16 +12,25 @@ ``test_start_worker__benchmark`` from ``test_process_queue__benchmark`` and divide the queue depth by the difference to get the marginal throughput of a busy queue. -Threadmill, django-tasks-db and django-tasks-redis run one worker process that -drains a queue and exits. Celery has no such mode, so the benchmark queues a -sentinel task last and waits for it to be processed. That wait is what proves the -queue was drained. Its worker is stopped after the measurement, because a graceful -shutdown takes seconds and would dominate a short drain. +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 to be processed. That wait is what +proves the queue was drained. Their workers are stopped after the measurement, +because a graceful shutdown takes seconds and would dominate a short drain. + +Celery and dramatiq read four messages per worker thread ahead. Dramatiq polls +instead of blocking and sleeps a jittered backoff once its window is full, so at +four messages its drain still measures that backoff more than its queue: the same +worker reaches thousands of tasks per second at a window deep enough to stay out +of the sleep. Threadmill, django-tasks-db and django-tasks-rq read one task at a +time, because their workers block on an empty queue and gain nothing from a +window. RQ's worker forks a work horse per job, so its drain includes that fork. """ import collections.abc import dataclasses import io +import os import subprocess import sys import tempfile @@ -46,6 +55,7 @@ celery_echo, celery_mark_processed, ) +from benchmarks.dramatiq_app import dramatiq_echo, dramatiq_mark_processed from tests.testapp.tasks import echo ENQUEUE_ITERATIONS = 500 @@ -54,6 +64,16 @@ QUEUE_DEPTH = 5000 """Tasks queued before one processing benchmark round.""" +READ_AHEAD = 4 +"""Messages each worker reads ahead, where its queue has such a setting. + +Celery ships a prefetch multiplier of four. Dramatiq's own default is two per +worker thread, and neither is deep enough for a single-threaded dramatiq to stop +measuring its poll backoff: the same worker drains 451 tasks/s at two messages, +1,983 at 16 and 5,658 at 64, where the backoff it sleeps once the window is full +stops setting the rate. +""" + CELERY_WORKER = ( sys.executable, "-m", @@ -62,16 +82,37 @@ "benchmarks.celery_app:celery_app", "worker", "--pool=solo", - "--prefetch-multiplier=1", + f"--prefetch-multiplier={READ_AHEAD}", "--loglevel=WARNING", "--without-gossip", "--without-mingle", "--without-heartbeat", ) -"""Celery worker running as one process with one thread, reading one message at a time. +"""Celery worker running as one process with one thread, ``READ_AHEAD`` messages ahead. The default prefork pool crashes on CPython 3.14, where the pool child loses -the task handler state it expects. +the task handler state it expects, so the worker runs on the solo pool. With one +concurrent task, ``--prefetch-multiplier`` sets the prefetch count to +``READ_AHEAD``, the benchmark rate. Its consumer blocks while the queue is empty, +so the window costs no sleep per message. +""" + +DRAMATIQ_WORKER = ( + sys.executable, + "-m", + "dramatiq", + "benchmarks.dramatiq_app:redis_broker", + "--processes", + "1", + "--threads", + "1", +) +"""dramatiq worker running as one process with one thread, ``READ_AHEAD`` messages ahead. + +The CLI has no read-ahead flag, so the worker environment carries +``dramatiq_queue_prefetch``. Its Redis consumer polls rather than blocks and +sleeps a jittered backoff once its window fills, so ``READ_AHEAD`` still bounds +this drain; see that constant for what a deeper window measures. """ WORKER_STOP_TIMEOUT_SECONDS = 20 @@ -134,11 +175,13 @@ def drain_with_django_tasks_db_worker() -> None: ) -def drain_with_django_tasks_redis_worker() -> None: - """Process every queued task with the django-tasks-redis worker.""" +def drain_with_django_tasks_rq_worker() -> None: + """Process every queued task with the django-tasks-rq worker.""" call_command( - "run_redis_tasks", - backend_name="django-tasks-redis", + "rqworker", + "--burst", + "--job-class", + "django_tasks_rq.Job", verbosity=0, stdout=io.StringIO(), ) @@ -150,7 +193,19 @@ def drain_with_celery_worker() -> None: drain_with_subprocess_worker(CELERY_WORKER) -def drain_with_subprocess_worker(argv: collections.abc.Sequence[str]) -> None: +def drain_with_dramatiq_worker() -> None: + """Process every queued task with a single-process, single-thread dramatiq worker.""" + dramatiq_mark_processed.send() + drain_with_subprocess_worker( + DRAMATIQ_WORKER, + env={**os.environ, "dramatiq_queue_prefetch": str(READ_AHEAD)}, + ) + + +def drain_with_subprocess_worker( + argv: collections.abc.Sequence[str], + env: collections.abc.Mapping[str, str] | None = None, +) -> None: """Run a worker CLI until the sentinel task queued last was processed.""" client = redis.Redis.from_url(REDIS_URL) client.delete(PROCESSED_KEY) @@ -158,6 +213,7 @@ def drain_with_subprocess_worker(argv: collections.abc.Sequence[str]) -> None: # The command is a fixed worker CLI, never caller input. process = subprocess.Popen( # noqa: S603 argv, + env=env, stdout=log, stderr=subprocess.STDOUT, ) @@ -220,10 +276,15 @@ def enqueue(count: int) -> TaskResult: return [task.enqueue(index) for index in range(count)][-1] def verify_processed(enqueued_task_result: TaskResult | None) -> None: - """Assert that the backend executed the benchmark tasks.""" + """Assert that the backend executed the benchmark tasks. + + django-tasks-rq returns the ``django-tasks`` backport's own + ``TaskResultStatus``, a different enum class carrying the same string + value, so the status is compared by value rather than identity. + """ assert enqueued_task_result is not None, "enqueue() must return a task result" task_result = task_backends[alias].get_result(enqueued_task_result.id) - assert task_result.status is TaskResultStatus.SUCCESSFUL, ( + assert task_result.status == TaskResultStatus.SUCCESSFUL, ( f"{alias} did not execute the benchmark tasks" ) @@ -236,6 +297,12 @@ def enqueue_celery_tasks(count: int) -> None: celery_echo.delay(index) +def enqueue_dramatiq_tasks(count: int) -> None: + """Accept `count` echo tasks on the dramatiq queue.""" + for index in range(count): + dramatiq_echo.send(index) + + def django_task_backend( name: str, alias: str, @@ -254,13 +321,18 @@ def django_task_backend( "django-tasks-db", "django-tasks-db", drain_with_django_tasks_db_worker ), django_task_backend( - "django-tasks-redis", "django-tasks-redis", drain_with_django_tasks_redis_worker + "django-tasks-rq", "django-tasks-rq", drain_with_django_tasks_rq_worker ), QueueUnderTest( name="celery", enqueue=enqueue_celery_tasks, drain=drain_with_celery_worker, ), + QueueUnderTest( + name="dramatiq", + enqueue=enqueue_dramatiq_tasks, + drain=drain_with_dramatiq_worker, + ), ) """Queues that ship a worker to process queued tasks.""" @@ -297,11 +369,18 @@ def stop_workers(empty_queues): @pytest.fixture def empty_queues(): - """Delete queued tasks from every compared queue before and after a benchmark.""" + """Delete queued tasks and stored results from every compared queue before and after a benchmark.""" client = task_backends[DEFAULT_TASK_BACKEND_ALIAS].client def delete_queued_tasks() -> None: - for key_pattern in ("threadmill:*", "django_tasks:*", "celery*", "_kombu*"): + for key_pattern in ( + "threadmill:*", + "django_tasks:*", + "celery*", + "dramatiq:*", + "rq:*", + "_kombu*", + ): if keys := client.keys(key_pattern): client.delete(*keys) client.delete(PROCESSED_KEY) diff --git a/docs/images/backend-comparison-dark.svg b/docs/images/backend-comparison-dark.svg index c58fd71..488f20a 100644 --- a/docs/images/backend-comparison-dark.svg +++ b/docs/images/backend-comparison-dark.svg @@ -1,19 +1,22 @@ - + - + Queue throughput 5,000 trivial tasks per queue · one worker process, one thread · higher is better threadmill -5,069/s +5,023/s celery - -2,128/s + +1,951/s django-tasks-db - -1,997/s -django-tasks-redis - -1,313/s -One message in flight per worker — no queue reads ahead. + +1,942/s +dramatiq + +669/s +django-tasks-rq + +80/s +threadmill, django-tasks-db and django-tasks-rq read one message at a time; celery and dramatiq four. diff --git a/docs/images/backend-comparison-light.svg b/docs/images/backend-comparison-light.svg index 3d94c84..cca0b6e 100644 --- a/docs/images/backend-comparison-light.svg +++ b/docs/images/backend-comparison-light.svg @@ -1,19 +1,22 @@ - + - + Queue throughput 5,000 trivial tasks per queue · one worker process, one thread · higher is better threadmill -5,069/s +5,023/s celery - -2,128/s + +1,951/s django-tasks-db - -1,997/s -django-tasks-redis - -1,313/s -One message in flight per worker — no queue reads ahead. + +1,942/s +dramatiq + +669/s +django-tasks-rq + +80/s +threadmill, django-tasks-db and django-tasks-rq read one message at a time; celery and dramatiq four. diff --git a/pyproject.toml b/pyproject.toml index 340e138..be1fd9d 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -51,7 +51,8 @@ dev = [ test = [ "celery", "django-tasks-db", - "django-tasks-redis", + "django-tasks-rq", + "dramatiq[redis]", "pytest", "pytest-asyncio", "pytest-benchmark", diff --git a/tests/testapp/settings.py b/tests/testapp/settings.py index 4522db2..c315cbc 100644 --- a/tests/testapp/settings.py +++ b/tests/testapp/settings.py @@ -43,8 +43,9 @@ "django.contrib.staticfiles", "threadmill", "tests.testapp", + "django_rq", "django_tasks_db", - "django_tasks_redis", + "django_tasks_rq", ] MIDDLEWARE = [ @@ -104,10 +105,9 @@ "BACKEND": "django_tasks_db.DatabaseBackend", "QUEUES": [DEFAULT_TASK_QUEUE_NAME], }, - "django-tasks-redis": { - "BACKEND": "django_tasks_redis.RedisTaskBackend", + "django-tasks-rq": { + "BACKEND": "django_tasks_rq.RQBackend", "QUEUES": [DEFAULT_TASK_QUEUE_NAME], - "OPTIONS": {"REDIS_URL": REDIS_URL}, }, "immediate": { "BACKEND": "django.tasks.backends.immediate.ImmediateBackend", @@ -117,6 +117,11 @@ }, } +# django-rq resolves every queue named in TASKS from this mapping. +RQ_QUEUES = { + DEFAULT_TASK_QUEUE_NAME: {"URL": REDIS_URL}, +} + # Run workers as quietly as the celery and dramatiq benchmarks. Django's # default logging pins django.tasks to INFO, so quiet it explicitly; the # test app's task logs stay visible.