Balance the transfer load between the workers - #1258
Open
gregorjerse wants to merge 3 commits into
Open
gregorjerse wants to merge 3 commits into
gregorjerse wants to merge 3 commits into
Conversation
The terminate handler cancelled every task, but the transfer ran in a thread pool whose context manager joined the workers on exit. The event loop blocked until all downloads finished: the heartbeat stopped and the container kept downloading data nobody needed. The bare except in transfer_inputs also reported the cancellation as a transfer error. Transfer.abort() stops the workers before the next object and closes the streams of the objects in transfer, so the running downloads fail at once; the aborted objects are removed and not retried. The executor is no longer joined on cancel, the terminate handler cancels only the main task, and the shared storage download reports DOWNLOAD_ABORTED so the location is released.
This was referenced Sep 22, 2026
A transient error failed the init container after three attempts spread over three seconds. The retry list named requests exceptions boto3 no longer raises and missed the botocore connection and read errors, the failure surfaced only after every worker drained its chunk, the partial object was hashed in full before the retry, and the error listed every file without naming the failed one. The transfer of an object is now attempted eight times with 1, 2, ..., 32, 60 seconds between the attempts; the wait ends on abort. The botocore connection, HTTP client and incomplete read errors are retried. The partial object is removed before the retry. The first failed worker aborts the others and its failure, naming the object, is reported.
The files were split into one fixed slice per worker, so a worker that drew the large files ran alone long after the others finished, and the init container used a single worker for less than 200 files. The workers now share one iterator over the files sorted largest first and each takes the next file when it finishes the previous one. The init container always uses GENESIS_MAX_DOWNLOAD_THREADS workers (default 3), capped by the number of files.
gregorjerse
force-pushed
the
fix-transfer-parallelism
branch
from
September 22, 2026 06:51
8a923b0 to
ae4ef46
Compare
This was referenced Sep 22, 2026
dblenkus
reviewed
Sep 22, 2026
dblenkus
left a comment
Member
There was a problem hiding this comment.
Best of the stack — the shared queue is the right fix and the diff is smaller than what it replaces. Dropping chunks also removes an accidental step-slice of a QuerySet in verify_data. Reviewing the last commit only.
paralelizenow always returns already-done futures (thewithjoins on the way out). The docstring says so, which is honest, but the "collection of futures, caller checks them" shape is vestigial — both callers immediately.result().list[bool]would be simpler while the function is being rewritten anyway.verify_datagains the shared queue but not largest-first —largest_firstis applied only intransfer_objects.ReferencedPathhas a size, so it could be ordered too.paralelize/async_paralelizedoobjects = list(objects)on the listlargest_firstjust built.
One interaction worth noting: the shared queue means a failing transfer_objects now retries every object rather than one worker's slice, which multiplies the retry-budget concern I raised on #1257.
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.
Stacked on #1257 and #1256. Until they merge, review the last commit only.
Problem
The transfer helpers split the files into one fixed slice per worker (
files[i::n]), regardless of size. A worker that drew the large files ran alone long after the others were idle, so the download took as long as the unluckiest slice. The init container also derived the worker count from the file count, one per 100 files, so any input set under 200 files was downloaded by a single worker.Fix
SharedIterator,largest_first). Starting the big files first and assigning dynamically keeps every worker busy until the queue drains, so they finish at about the same time, also when a file transfers slower than its size suggests. This is longest-first list scheduling. It lives inparalelizeandasync_paralelize, soTransfer.transfer_objects(uploads, shared-storage downloads) andStorageLocation.verify_datagain the same balance.GENESIS_MAX_DOWNLOAD_THREADSworkers (default 3), capped by the number of files.asyncio.wait.chunksis removed; nothing else used it.Considered and not done
A static split of the bytes into equal shares per worker balances only on paper: the throughput of a file depends on its size (small files are latency bound), on the S3 side and on the moment, so fixed shares still finish apart. The shared queue balances on actual progress. One file larger than everything else still bounds the download by its own transfer time; that is set by boto's per-file concurrency (10 ranged GETs) and the volume throughput, not by scheduling.
Verification
largest_firstordering; a balancing test with one slow and nine fast items on two workers, where the old static split gives 5 and 5 items per worker and the shared queue 1 and 9.resolwe.storage.tests.test_utils,test_transfer,test_modelsandresolwe.flow.tests.test_transferpass (50 tests).