Skip to content

Balance the transfer load between the workers - #1258

Open
gregorjerse wants to merge 3 commits into
genialis:masterfrom
gregorjerse:fix-transfer-parallelism
Open

gregorjerse wants to merge 3 commits into
genialis:masterfrom
gregorjerse:fix-transfer-parallelism

Conversation

@gregorjerse

Copy link
Copy Markdown
Contributor

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

  • The workers share one thread-safe iterator over the files sorted largest first and each takes the next file when it finishes the previous one (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 in paralelize and async_paralelize, so Transfer.transfer_objects (uploads, shared-storage downloads) and StorageLocation.verify_data gain the same balance.
  • The init container always uses GENESIS_MAX_DOWNLOAD_THREADS workers (default 3), capped by the number of files.
  • Empty input returns no workers instead of raising from asyncio.wait.
  • chunks is 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

  • New tests: the shared iterator hands out every item exactly once across four threads; largest_first ordering; 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_models and resolwe.flow.tests.test_transfer pass (50 tests).

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

@dblenkus dblenkus left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

  1. paralelize now always returns already-done futures (the with joins 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.
  2. verify_data gains the shared queue but not largest-first — largest_first is applied only in transfer_objects. ReferencedPath has a size, so it could be ordered too.
  3. paralelize/async_paralelize do objects = list(objects) on the list largest_first just 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.

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.

2 participants