Skip to content

Prefetch tasks in per-process batches - #65

Merged
codingjoe merged 27 commits into
mainfrom
codingjoe-task-prefetching
Oct 7, 2026
Merged

codingjoe merged 27 commits into
mainfrom
codingjoe-task-prefetching

Conversation

@codingjoe

@codingjoe codingjoe commented Oct 1, 2026 •

Copy link
Copy Markdown
Owner

Each task cost two Redis round-trips on the worker's critical path, so for fast tasks a worker thread spent more wall-clock waiting on the broker than running work.

  • acquire grows a count and returns a list, so a worker reserves up to count tasks in one atomic EVAL; acquire.lua pops them round-robin across queues.
  • A per-process fetcher thread fills a bounded buffer that worker threads drain, sized by --prefetch-count, which defaults to four times the thread count and accepts 1 to disable batching.
  • The buffer is a priority queue: it dispatches the highest priority task first and keeps queue order within a priority, so a batch that interleaves queues still runs by priority.
  • Graceful shutdown drains the buffer before the process exits, and a hard kill is left to the lease reaper.
  • A fetch failure is logged through the structured pipeline and exits the child non-zero rather than reading as a clean drain, and a worker whose consumer threads all died is recycled instead of parking.
  • Merges main: the lease lives beside the payload with its own token per task, so each prefetched task carries its lease and an acknowledgement is dropped once a later attempt holds the lease.
  • The benchmark measures dramatiq and django-tasks-rq beside the existing queues, each reading 128 messages ahead where its queue allows it, with per-queue depths and a guard that fails a degenerate measurement instead of charting noise.

Measured on the merged code with one worker process and one thread: threadmill 11,973 tasks/s, dramatiq 7,331, celery 2,307, django-tasks-db 2,080, django-tasks-rq 90. Against a local broker the prefetch buffer is worth about a fifth over reading one at a time, 11,973 against 9,972.

Trade-offs: tasks are leased when they are fetched, so time in the buffer counts against lease_ttl; a task enqueued after a fetch waits for the buffer to drain before a worker picks it up; and --max-tasks may overshoot by up to the buffer because a prefetched task always runs.

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.
…ching

# Conflicts:
#	README.md
#	benchmarks/chart.py
#	benchmarks/test_backends.py
#	docs/images/backend-comparison-dark.svg
#	docs/images/backend-comparison-light.svg
#	pyproject.toml
#	tests/test_executor.py
#	threadmill/backends/lua/acquire.lua
#	threadmill/backends/redis.py
The merged worker leases tasks beside their payload instead of rewriting
it, which roughly doubled threadmill's drain: 11,973 tasks/s against
5,437 before the merge, and it now leads the comparison. The prefetch
buffer is worth about a fifth over reading one at a time, 11,973 against
9,972, where the same merge left the ablation at its pre-merge rate.

The chart now carries django-tasks-rq, which main added in place of
django-tasks-redis, and drops the dependency the replacement orphaned.
Its queue drains at 90 tasks/s, so it measures 5,000 tasks.
The buffer hands out the highest priority task first and keeps fetch order
within a priority, so prefetching weakens the queue's ordering less: a
batch that interleaves queues still dispatches by priority.

It holds a small ordered wrapper rather than bare task results, keyed on
the negated priority, because the backend pops the highest priority first,
and a fetch sequence counter that keeps equal priorities in queue order.
A task waits in the buffer while it holds its lease, so a deep buffer needs
a matching lease_ttl. That is a downside the operator sizes with
--prefetch-count, not something the worker should compensate for, so the
note that suggested renewing leases is gone and the trade-off is stated
where the knob is set: the README's soft-limit list and the flag's help.
The wait that bounds a broker acquire and a buffer get moves from a module
global onto WorkerProcess, which already owns its threads' configuration
and is carried by both the fetcher and the consumers.

The prefetch count now derives in one line: a falsy value or zero takes
four tasks per thread, an explicit value is kept, and a negative value is
floored at one, because a negative count left a worker that fetched
nothing while the pool respawned it. The CLI still rejects a negative
count outright, so an operator gets an error rather than a silently
adjusted configuration.
It repeated counts the operator had already configured, and the command
prints that it is starting workers. The test that watched the line now
asserts on the executor the command builds, the way the poll-interval
tests do, so the prefetch plumbing stays pinned without a log to watch.
The prefetch budget clamp only shrank the last request, so it bought a
margin of at most one task on a limit the README already called soft. The
worker's expiry, set when the budget is reached and read at the top of the
fetch loop, already stops fetching, so the clamp and its remaining_tasks
helper are gone and the fetcher always asks for a full buffer.

The documented overshoot had to grow with it: the buffer and the batch in
the fetcher's hand both still run, so a recycled worker can finish roughly
twice the buffer size plus the thread count past its budget, which is what
the README and the flag help now say.
ThreadmillTaskResult gains a comparison ordered by descending priority and
then enqueue time, which is the order the queue pops in, so the prefetch
buffer holds the leased results directly. That removes the PrefetchedTask
wrapper, the fetch-sequence counter and the unwrapping at the consumer.

The acquire contract now says what the buffer needs: it returns
ThreadmillTaskResult, which carries the lease and is orderable.
Short sentences, one fact for each sentence, and no em dashes, semicolons
or colons in the prose. The acquire docstring says what the method locks
and what it returns. The prefetcher and worker docstrings say what each
thread does instead of naming it again.

The range notation in the benchmark docs becomes words, "5 to 10 ms", and
the definition list of benchmarks becomes three sentences, because a colon
after each name reads as structure and not as prose.
The section said the same facts in more words. It now uses short sentences
and one idea for each sentence, and it drops the em dashes, the semicolons
and the two colons that joined clauses. The prefetch buffer paragraph, the
health note, the shutdown note and the lease note all lost a third of
their length and kept every number.
The method is a dunder, and the class does not document its other dunders
either. The expression states the order on its own.
The try block in run now holds one call, so the handler shows what it does.
It stores the failure for the worker process, logs it and raises it again.
The loop that fills the buffer lives in fill_buffer.
The fetcher had two signals for one event, an event for completion and a
field for the exception, and a failure was reported three times. It reached
the structured log, came out again through the default thread hook and came
out a third time when the worker raised it.

One Future now carries the outcome. The fetch loop sets an exception on it
or a result, the consumers wait for done, and the worker reads it once after
the join. A failure is logged once with its traceback, and the child exits
with code 1 through SystemExit, which prints nothing of its own.
The section explained the same feature three times over, once as prose,
once as the flag and once as four soft limits. The flag help carries the
size and the default, the Redis options carry the lease effect, and the
health note carries the max-tasks overshoot, so the README keeps one
paragraph that says what the extra thread does and which flag sizes it.
The buffer is a plain queue again, and the comparison on
ThreadmillTaskResult is gone with it. A worker takes the task that the
fetcher leased first, which is the order the queue handed them out.

The ordering rule is unchanged from before prefetching. The backend still
picks each queue's head by priority and enqueue time, and the buffer now
preserves that lease order instead of reordering it.
The flag name and its default are enough for the help. The lease effect
lives with the Redis options and the max-tasks overshoot lives in the
health note, so the help no longer repeats them.
The future replaced a finished event that a finally block always set, so
the new code lost that guarantee. A BaseException escaping the fetch loop
left the future unresolved, and a consumer waits for a result that never
comes. The finally block now resolves the future when the handler has not,
and the BaseException still escapes the thread.
The guard reads on its own.
The flag help documents the setting, so the README does not repeat it.
A forkserver child inherits the stdout of the long-lived forkserver
instead of the file descriptor capfd replaces, so the fetch failure
report never reached the fixture and both tests failed on Linux, where
Python 3.14 defaults to forkserver. Force spawn for the duration of the
child, like test_run__routes_task_logs_to_stdout does.
Regenerated from a run on the current branch. Threadmill holds 11,977
tasks per second, dramatiq 7,168, celery 2,183, django-tasks-db 2,154 and
django-tasks-rq 87, which is within a few percent of the previous chart
for every row.

The no-prefetch ablation moved by a quarter, from the one-second
quantization of the worker start and stop rather than a change in the
drain. Its start mean fell by a second while its process mean rose by the
same second, so the row is noisy at this depth and the chart does not
plot it.
@codingjoe
codingjoe marked this pull request as ready for review October 7, 2026 21:01
@codingjoe
codingjoe merged commit 655f8a6 into main Oct 7, 2026
4 checks passed
@codingjoe
codingjoe deleted the codingjoe-task-prefetching branch October 7, 2026 21:02
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