Skip to content

Prefetch tasks in per-process batches - #65

Draft
codingjoe wants to merge 6 commits into
mainfrom
codingjoe-task-prefetching
Draft

codingjoe wants to merge 6 commits into
mainfrom
codingjoe-task-prefetching

Conversation

@codingjoe

@codingjoe codingjoe commented Oct 1, 2026 •

Copy link
Copy Markdown
Owner

Each task currently costs two Redis round-trips (acquire then acknowledge) on the critical path, so for fast tasks a worker thread spends most of its wall-clock waiting on the broker instead of running work. This patch reserves tasks ahead of time in one batched call, which amortizes that latency and keeps the pool busy, in the spirit of the utilization principle in CONTRIBUTING.md.

Approach

One acquire method, extended with a count, feeding one buffer per worker process:

  • ThreadmillTaskBackend.acquire(*queues, count=1, timeout=None, worker="") -> list[TaskResult] keeps its name and grows a count, so there is no second entry point to learn. It waits for the first task only, fills the rest without waiting, and never returns an empty list. Like the rest of the interface class, it declares the contract and raises NotImplementedError.
  • acquire.lua pops up to count tasks round-robin across queues and returns an array, so a batch remains a single atomic EVAL. acquire with count=1 behaves exactly as before.
  • A TaskPrefetcher daemon thread per worker process fills a bounded queue.Queue; worker threads drain it. A full buffer backpressures the fetcher, so memory is bounded by --prefetch-count, which defaults to 4 x threads and accepts 1 to disable batching.
  • Graceful shutdown finishes the buffered tasks (consumers drain before the fetcher is stopped). A hard kill is left to the lease reaper, which already covers it.
flowchart LR
    Q[(Ready queue)] -->|acquire, one round-trip for N tasks| F[TaskPrefetcher]
    F -->|bounded put| B[/task_buffer/]
    B --> W[WorkerThreads]
    W -->|acknowledge| R[(Results)]
Loading

Failure and lifecycle behaviour

A fetch error is logged through the structured pipeline, recorded on the fetcher, and re-raised so the child exits non-zero, rather than disappearing into a daemon thread as if the queue were drained. A worker whose consumer threads all died now recycles instead of parking in join forever, which the design required a dedicated stop signal for.

Benchmark chart

dramatiq joins the comparison. Every queue runs one process and one thread, and every queue that can be told reads 128 messages ahead, the rate an earlier harness in this repository used when dramatiq led the field:

queue read-ahead tasks/s
dramatiq 128 (dramatiq_queue_prefetch) 6,975
threadmill 128 (prefetch_count) 5,437
celery 128 (--prefetch-multiplier) 2,080
django-tasks-db 1, no setting available 1,981
django-tasks-redis 1, no setting available 1,391

Methodology, because the ranking is very sensitive to it and I had it wrong twice:

  • A shallow window measures the poll, not the queue. dramatiq's Redis consumer never blocks: with its window full it sleeps compute_backoff(0), a jittered 5-10 ms, so the cost per task is inverse in the read-ahead. Measured curve: 1 -> 111 tasks/s, 2 -> 220, 4 -> 434, 8 -> 838, 32 -> 2,866, 128 -> 6,975. An earlier revision of this branch pinned it to one message in flight and published 105 tasks/s; at its own default of two it read 211.
  • Result storage is matched. dramatiq takes its Results middleware (store_results=True) so it stores results the way celery does. The earlier harness did the same, and without it dramatiq is measured while discarding the value.
  • Per-queue depth. A cold worker start and stop is quantized to about a second, which swamped the marginal drain of a shallow queue and once produced a negative throughput. The threadmill queues run 60,000 tasks and the third-party queues 20,000, and a degenerate drain now fails the benchmark instead of writing a negative value.
  • The two Django backends read one at a time, because their shipped workers expose no read-ahead setting: db_worker claims one task per loop and run_redis_tasks hardcodes max_messages=1. Patching a third-party worker would measure the patch rather than the library, so they stay as they ship and the chart footnote says so.
  • Threadmill's own ablation: 5,437 tasks/s with batching against 4,994 reading one at a time, so the buffer is worth about nine percent even against a local broker. That row stays in the suite and out of the chart.

Why dramatiq leads

Measured per task rather than inferred, because my first explanation was wrong. The two queues issue about the same number of Redis commands - 11.03 against 11.15, the latter including dramatiq's five-command result store - but threadmill's commands are far more expensive:

per task threadmill dramatiq
Redis commands 11.03 11.15
server CPU 54.3 µs 11.7 µs
of which the script EVAL 43.4 µs 9.5 µs
server CPU, 5 KB payload 79.9 µs 21.2 µs

Threadmill makes two full JSON passes per task: the fetch script decodes the payload, stamps RUNNING with the lease and worker id, and re-encodes it, and the ack writes the whole serialized TaskResult plus a history entry, an eviction scan and a telemetry publish. dramatiq never rewrites a payload - its fetch is LPOP plus SADD per id - so its per-task server work is trivial. Redis is single-threaded, so that CPU is time every other client waits behind, which is the honest price of the lease visibility, durable results and telemetry threadmill offers and dramatiq does not. Filed as a separate optimization issue, since it is a redesign of the scripts rather than part of this branch.

Trade-offs worth reviewing

  • Buffer dwell counts against lease_ttl. Reserved tasks are marked RUNNING at fetch time, so a long queue of work ahead of a buffered task can outlive its lease and be reaped FAILED before it runs. This is documented in the README; sizing guidance is included. Renewing leases is the upgrade path if dwell ever matters.
  • Priority lookahead widens from threads to the buffer size, so ordering is no longer strictly global.
  • --max-tasks is a soft limit and can overshoot by up to the buffer, because a prefetched task is always finished.
  • worker_ids records the process-level fetcher identity, not the executing thread, since the fetcher holds the lease.

Deferred, filed as follow-ups

  • A ceiling for --prefetch-count (operator-controlled today; an unbounded value lets one worker take the whole backlog and blocks Redis for work proportional to the count).
  • Releasing or requeueing the leased batch a failed fetch orphans (the error path is loud now; returning the tasks is the larger design piece).
  • The per-task Redis cost above: stamping leased state outside the payload, batching the telemetry publish, and amortizing the history and eviction work.
  • Four untested pre-existing branches (telemetry, backend error raises, CLI error paths).

Testing

  • uv run pytest -m "not benchmark": 169 passed, 5 skipped (textual absent locally); tests/test_inspector.py with the extra: 58 passed.
  • The comparison benchmark: 20 passed in 1m 44s, every queue inside its timeout, all six rows measurable.
  • 100% patch coverage, line and branch (--cov-branch reports no partial branch on any added line).
  • uvx prek run --all-files green, and the commit hooks pass.

Add a count to ThreadmillTaskBackend.acquire so a worker reserves up to
`count` tasks in one broker round-trip, and fill a per-process buffer
from a dedicated fetcher thread. The buffer defaults to 4 x threads and
is tunable with --prefetch-count; 1 disables batching.

- Redis acquire pops a round-robin batch atomically and advances the
  rotation one position per call
- the fetcher is a daemon thread, stops on max_tasks, shutdown, or drain,
  abandons a full buffer once no consumer is left, and logs plus re-raises
  a fetch failure so the child exits non-zero
- a worker whose consumers died is recycled instead of parking forever
- document the option and its soft limits in the README
Measure dramatiq beside celery and the task backends on the same trivial
echo task, pinned to one process, one worker thread and a prefetch of one
message so it matches the others.

The threadmill queues grew to 60,000 tasks and dramatiq keeps 5,000: the
fixed cost of a cold worker start and stop is quantized to about a second,
which swamped the marginal drain of a shallower queue and made the
prefetch comparison unmeasurable.

- per-queue depth on QueueUnderTest, recorded in the benchmark extra info
- fail loudly instead of writing a negative throughput when a drain is
  degenerate (process mean below start mean)
- chart height follows the row count, and its subtitle reports the depths
  actually measured
- the chart plots threadmill at its default configuration only; the
  no-prefetch run stays a benchmark diagnostic, since on a local broker
  the two land within a percent of each other
dramatiq's Redis consumer polls rather than blocks: with its read-ahead
window full it sleeps compute_backoff(0), a jittered 5-10 ms, so pinning
it to one message in flight cost a sleep between every task and made the
chart read 105/s instead of its real figure. Unpinned it reads 211/s,
which is five times celery rather than twenty.

Celery's pinned prefetch multiplier goes with it: its consumer blocks on
Redis, so the pin measured nothing (2,145/s pinned against 2,080/s at its
default over 20,000 tasks).

The methodology is now one worker process and one thread, each queue at
its own default read-ahead, and the chart says so.
Threadmill reserves four tasks per worker by default, so every queue that
can be told now reads four ahead: celery through --prefetch-multiplier=4
and dramatiq through dramatiq_queue_prefetch=4. The prefetch buffer is no
longer a comparison advantage.

dramatiq reads 420 tasks/s at that rate, twice the 211 it scored at its
own default of two; its polling consumer still pays a jittered 5-10 ms
backoff roughly once per four messages, which the chart footnote states.

The two Django backends keep reading one message at a time because their
shipped workers expose no read-ahead setting: db_worker claims one task
per loop and run_redis_tasks hardcodes max_messages=1. Patching a third
party worker would measure the patch rather than the library, so they stay
as they ship and the footnote says so.
An earlier harness in this repository ran dramatiq with a prefetch window
of 128 and it led the field; a later commit dropped it because its
single-threaded consumer sleeps a poll backoff between messages and could
not be compared fairly against queues that read one message at a time.

That window is the fix, not the problem. A shared rate of four still left
dramatiq penalised: its jittered 5-10 ms backoff lands once per window, so
the cost per task is inverse in the depth and a shallow window measures the
poll, not the queue. Every queue that can be told now reads 128 ahead,
which amortises that backoff to about 0.06 ms per task and leaves celery
and threadmill blocking on Redis as they always did.

dramatiq also takes its Results middleware back, which stores results the
way celery does, so neither queue is measured discarding the value. Its
depth rises to 20,000 like the other third-party queues, because at this
rate a 5,000 task drain fits inside the one-second quantisation of the
fixed start cost.

dramatiq leads at 6,975 tasks/s against threadmill's 5,437, close to the
7,676 the earlier harness measured, and threadmill's own prefetch buffer
is worth about nine percent over reading one at a time.
Storing results for parity with celery left keys the cleanup could not
match: with the default result backend the key is a bare md5 hex, so the
harness's dramatiq:* pattern missed it and the results stayed in Redis
until their TTL expired.

Naming the result namespace makes the key greppable and the existing
pattern deletes it, so no new cleanup code is needed.
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant