Repository navigation
Prefetch tasks in per-process batches - #65
Merged
Merged
Conversation
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.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
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.
acquiregrows acountand returns a list, so a worker reserves up tocounttasks in one atomicEVAL;acquire.luapops them round-robin across queues.--prefetch-count, which defaults to four times the thread count and accepts1to disable batching.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-tasksmay overshoot by up to the buffer because a prefetched task always runs.