Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,7 @@
<picture>
<source media="(prefers-color-scheme: dark)" srcset="https://github.com/codingjoe/threadmill/raw/main/docs/images/backend-comparison-dark.svg">
<source media="(prefers-color-scheme: light)" srcset="https://github.com/codingjoe/threadmill/raw/main/docs/images/backend-comparison-light.svg">
<img alt="Tasks per second with one worker: threadmill 5,069, celery 2,128, django-tasks-db 1,997, django-tasks-redis 1,313." src="https://github.com/codingjoe/threadmill/raw/main/docs/images/backend-comparison-light.svg">
<img alt="Tasks per second with one worker: threadmill 5,023, celery 1,951, django-tasks-db 1,942, dramatiq 669, django-tasks-rq 80." src="https://github.com/codingjoe/threadmill/raw/main/docs/images/backend-comparison-light.svg">
</picture>
</p>

Expand Down
23 changes: 16 additions & 7 deletions benchmarks/chart.py
Original file line number Diff line number Diff line change
Expand Up @@ -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'
)
Expand Down Expand Up @@ -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'<svg xmlns="http://www.w3.org/2000/svg" viewBox="0 0 {WIDTH} {HEIGHT}" '
f'width="{WIDTH}" height="{HEIGHT}" role="img" '
f'<svg xmlns="http://www.w3.org/2000/svg" viewBox="0 0 {WIDTH} {height}" '
f'width="{WIDTH}" height="{height}" role="img" '
f'aria-label="{describe(results)}.">',
"<style>svg{max-width:100%;height:auto}</style>",
f'<rect x="0.5" y="0.5" width="{WIDTH - 1}" height="{HEIGHT - 1}" rx="14" '
f'<rect x="0.5" y="0.5" width="{WIDTH - 1}" height="{height - 1}" rx="14" '
f'fill="{theme.canvas}" stroke="{theme.border}"/>',
text(28, 46, "Queue throughput", theme=theme, size=19, weight=700),
text(
Expand Down Expand Up @@ -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,
Expand Down
49 changes: 49 additions & 0 deletions benchmarks/dramatiq_app.py
Original file line number Diff line number Diff line change
@@ -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:<queue>:<actor>:<message_id>`` 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)
115 changes: 97 additions & 18 deletions benchmarks/test_backends.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand All @@ -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",
Expand All @@ -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
Expand Down Expand Up @@ -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(),
)
Expand All @@ -150,14 +193,27 @@ 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)
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,
)
Expand Down Expand Up @@ -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"
)

Expand All @@ -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,
Expand All @@ -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."""

Expand Down Expand Up @@ -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)
Expand Down
25 changes: 14 additions & 11 deletions docs/images/backend-comparison-dark.svg
Loading
Sorry, something went wrong. Reload?
Sorry, we cannot display this file.
Sorry, this file is invalid so it cannot be displayed.
Loading
Loading