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 11,977, dramatiq 7,168, celery 2,183, django-tasks-db 2,154, django-tasks-rq 87." src="https://github.com/codingjoe/threadmill/raw/main/docs/images/backend-comparison-light.svg">
<img alt="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." src="https://github.com/codingjoe/threadmill/raw/main/docs/images/backend-comparison-light.svg">
</picture>
</p>

Expand Down
4 changes: 2 additions & 2 deletions benchmarks/chart.py
Original file line number Diff line number Diff line change
Expand Up @@ -221,10 +221,10 @@ def build_chart(results: list[QueueResult], theme: Theme) -> str:
text(
28,
footnote_y,
# joe: width checked by hand (right edge 856.1 of 900 at 11.5px); add a
# 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 and -rq read one message at a time.",
"128 ahead; django-tasks-db, -rq and huey read one task at a time.",
theme=theme,
size=11.5,
fill=theme.faint,
Expand Down
34 changes: 34 additions & 0 deletions benchmarks/huey_app.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
"""The huey app the comparison benchmarks hand tasks to.

The huey consumer CLI imports this module without setting up Django, so it must
not import Django or any Django application. The Redis storage keeps task
results in the ``huey.results.<name>`` hash, the way the Celery app's result
backend keeps them in Redis. The benchmark cleanup's ``huey.*`` pattern deletes
the queue, the results, the schedule and the counters.
"""

import os

import redis
from huey import RedisHuey

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)

huey_app = RedisHuey("threadmill_benchmark", results=True, url=REDIS_URL)


@huey_app.task()
def huey_echo(value):
"""Return the given value."""
return value


@huey_app.task()
def huey_mark_processed():
"""Record that every earlier task in the queue has been processed."""
client.incr(PROCESSED_KEY)
45 changes: 43 additions & 2 deletions benchmarks/test_backends.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,8 @@

django-tasks-db reads one task at a time, because its shipped worker does not
expose a read-ahead setting. django-tasks-rq forks a work horse for each job, so
its drain includes that fork and it reads one task at a time too.
its drain includes that fork and it reads one task at a time too. huey blocks on
an empty queue, but pops one message at a time either way.

Threadmill is measured twice, with 128 messages ahead and with one message at a
time. The read-ahead cost can therefore be subtracted from both worker benchmarks.
Expand Down Expand Up @@ -74,6 +75,7 @@
celery_mark_processed,
)
from benchmarks.dramatiq_app import dramatiq_echo, dramatiq_mark_processed
from benchmarks.huey_app import huey_echo, huey_mark_processed
from tests.testapp.tasks import echo

ENQUEUE_ITERATIONS = 500
Expand All @@ -96,7 +98,8 @@
about 0.06 ms for each task.

django-tasks-db reads one task at a time. django-tasks-rq forks a work horse for
each job. Neither queue can be told to read ahead.
each job. Neither queue can be told to read ahead. huey pops one task at a time
and has no read-ahead setting either.
"""

CELERY_WORKER = (
Expand Down Expand Up @@ -142,6 +145,26 @@
``dramatiq_queue_prefetch``. The benchmark sets this variable to ``READ_AHEAD``.
"""

HUEY_WORKER = (
sys.executable,
"-m",
"huey.bin.huey_consumer",
"benchmarks.huey_app.huey_app",
"--workers=1",
"--worker-type=thread",
"--no-periodic",
"--quiet",
"--graceful-signal=TERM",
)
"""huey consumer running as one process with one thread, 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
workers and ``--no-periodic`` skips the periodic-task scan, because the
benchmark app registers none. ``--graceful-signal=TERM`` stops the consumer on
SIGTERM the way the other worker benchmarks stop theirs.
"""

WORKER_STOP_TIMEOUT_SECONDS = 20
"""Seconds to wait for a worker process to stop after SIGTERM."""

Expand Down Expand Up @@ -251,6 +274,12 @@ def drain_with_dramatiq_worker() -> None:
)


def drain_with_huey_worker() -> None:
"""Process every queued task with a single-thread huey consumer."""
huey_mark_processed()
drain_with_subprocess_worker(HUEY_WORKER)


def drain_with_subprocess_worker(
argv: collections.abc.Sequence[str],
env: collections.abc.Mapping[str, str] | None = None,
Expand Down Expand Up @@ -352,6 +381,12 @@ def enqueue_dramatiq_tasks(count: int) -> None:
dramatiq_echo.send(index)


def enqueue_huey_tasks(count: int) -> None:
"""Accept `count` echo tasks on the huey queue."""
for index in range(count):
huey_echo(index)


def django_task_backend(
name: str,
alias: str,
Expand Down Expand Up @@ -405,6 +440,11 @@ def django_task_backend(
enqueue=enqueue_dramatiq_tasks,
drain=drain_with_dramatiq_worker,
),
QueueUnderTest(
name="huey",
enqueue=enqueue_huey_tasks,
drain=drain_with_huey_worker,
),
)
"""Queues that ship a worker to process queued tasks."""

Expand Down Expand Up @@ -457,6 +497,7 @@ def delete_queued_tasks() -> None:
"django_tasks:*",
"celery*",
"dramatiq:*", # broker keys and the dramatiq:results:* results
"huey.*", # queue, results, schedule and counter keys
"rq:*",
"_kombu*",
):
Expand Down
33 changes: 18 additions & 15 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.
33 changes: 18 additions & 15 deletions docs/images/backend-comparison-light.svg
Loading
Sorry, something went wrong. Reload?
Sorry, we cannot display this file.
Sorry, this file is invalid so it cannot be displayed.
1 change: 1 addition & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@ test = [
"django-tasks-db",
"django-tasks-rq",
"dramatiq[redis]>=2.2.1",
"huey",
"pytest",
"pytest-asyncio",
"pytest-benchmark",
Expand Down
Loading