Skip to content

feat(get): --workers streams to S3 across compute nodes (0.25.0) - #186

Merged
johnyaku merged 1 commit into
mainfrom
feat/get-workers
Aug 20, 2026
Merged

johnyaku merged 1 commit into
mainfrom
feat/get-workers

Conversation

@johnyaku

@johnyaku johnyaku commented Aug 20, 2026 •

Copy link
Copy Markdown
Contributor

Closes #172.

dt get had no distributed mode. dt push --workers had solved multi-node transfer already and the machinery was never push-specific, so this reuses it rather than rewriting it.

Scoped to s3://, deliberately

-o bound by a second node adds
s3://… the network bandwidth — a separate NIC
some/dir/ the filesystem both ends share contention

A local destination on Lustre is limited by the filesystem, not the node, so --workers is a hard usage error there, pointing at -j:

$ dt get myrepo data/x -o out/ -w 4
Error: --workers only applies to an s3:// destination.
Streaming to object storage is bound by the network, so more nodes mean more
bandwidth. A local destination is bound by the filesystem both ends share, where
a second node adds contention rather than throughput -- raise -j instead.

On NCI there is a second reason to like the S3 case: copyq is the queue with outbound network access, and it is already hpc.DEFAULT_TRANSFER_QUEUE.

Reuse, as the issue asked

  • hpc.partition_by_size() — the LPT bin-packer, lifted out of push.partition_manifest. What stays behind in push is the part that genuinely is push-specific: a file is named by its hash, so its size has to be read from the local cache. Splitting exactly there seemed better than a blind move.
  • hpc.save_manifest(..., metadata=...) — everything a worker needs beyond its own partition travels this way, because it is a fresh process on another machine. push needs a remote name; get also needs a destination, credentials config, and the flags the run was invoked with.
  • submit_workers / monitor_jobs / load_worker_partition — unchanged. The only accommodation needed was making dt get's REPOSITORY positional optional, since a worker reads everything from the manifest.

The docstring on partition_by_size records why LPT belongs to qxub workers and not to a ThreadPoolExecutor, so the next person doesn't apply it to a thread pool: separate processes cannot steal work, a shared queue self-corrects, and freezing an assignment the queue would fix is strictly worse.

The design questions, settled

Reachability — resolve once on the submitter and pass paths. The submitter has to resolve every row anyway (LPT needs byte sizes; list_source_files is where sizes come from), so having paid for that it ships the resolved md5, size and destination key per file. No worker ever runs dvc list — which is what keeps N compute nodes off the shared clone's SQLite state db. _run_dvc_list's retry loop exists because that lock is real. A worker opens the clone once, read-only, for one thing: the source remote's filesystem and credentials out of its .dvc/config.

Link types and the destination identity are resolved on the submitter too, not re-derived per worker. resolve_link_types() reads dvc config from the cwd and preflight() resolves a credential chain; both could answer differently on a compute node, and a transfer that silently changes its mind halfway is worse than one that fails.

Partition unit — per file, not per row. The issue framed this as a trade against per-row reporting, but it isn't one: each task carries its row tag, so reporting stays per-row and in CSV order even when a sample directory spans nodes. Per-file also balances properly in the two cases per-row cannot — a single row of thousands of files, and rows differing by an order of magnitude, which fastq samples reliably do.

Reporting — workers write result_<n>.json; the submitter collates by row tag. A row whose files do not all come back is reported failed even when every result that did arrive succeeded, because a job killed on walltime would otherwise read as a clean run of a short row.

Observed end-to-end

Real CLI, real partitioning, real manifest round-trip; fake S3 and fake qxub. Note row AF013-A was split across workers 0 and 1 and still reports as one row:

$ dt get bcarc_wts --csv samples.csv -o s3://collab-bucket/fastqs/ -w 3 -v
Destination: s3://collab-bucket/fastqs
Writing as:  account 533267394226 as arn:aws:iam::533267394226:user/collab
Partition byte balance across 3 active worker(s): min=600, max=3.91k, total=8.30k
✓ data/fq/AF013-A: 2 files -> s3://collab-bucket/fastqs/AF013-A
✓ data/fq/AF013-B: 1 file -> s3://collab-bucket/fastqs/AF013-B
✓ data/fq/AF013-C: 1 file -> s3://collab-bucket/fastqs/AF013-C
✗ data/fq/GONE: Not found in the source repository: data/fq/GONE

3 uploaded, 1 failed        # exit 1

worker_0.json: 1 file(s), 4000 bytes ['AF013-A/R1.fq']
worker_1.json: 1 file(s), 3900 bytes ['AF013-A/R2.fq']
worker_2.json: 2 file(s),  600 bytes ['AF013-B/R1.fq', 'AF013-C/R1.fq']

Security

Each worker preflights for itself, so --dest-account-id is enforced on the node that actually writes — a credential chain can resolve differently on a compute node. The manifest carries the profile name, endpoint and region; no key material. There's a test asserting that.

Also

The issue noted the ThreadPoolExecutor idiom copy-pasted five times, with "if a sixth copy is about to appear, factor it then". This adds a sixth use, so get.py's two are now behind _place_all and _upload_all, and _resolve_rows is shared by the local and S3 row-resolve paths. Side effect: upload_to_s3 now resolves rows concurrently, which it did not before. A cross-module utils.run_parallel() is still not worth it.

Still open — please read before merging

Whether more nodes actually move more bytes is not verified here, and needs a real submission. #172 asks for that measurement and I could not do it: it needs real data, a queue slot and time. It also isn't only about Lustre on the write side — if the source remote is on /g/data rather than object storage, the read side is a shared filesystem again and the same contention argument applies to it. docs/get.md says so, and says to time a single-node -j 8 run and watch for saturation before scaling up.

So: this ships the tool to do the measurement with, not the conclusion.

Tests

30 new unit tests. The central one replaces hpc.submit_workers with a stub that runs the worker in-process, so the manifest is exercised as a real interface rather than asserted against its spelling. Full unit suite: 2381 passed, 2 skipped, 0 failed — run on a compute node via qxub exec --env dt -- pytest tests/unit -q (2:03).

For the record, since an earlier version of this description said otherwise: a run on the login node reported 1 failure in test_auth_credentials_aws.py (FileNotFoundError: 'dvc'). That was an artefact of my test invocation, not a pre-existing failure and not this change — I had bypassed conda run, which leaves the env's bin off PATH, so anything shelling out to dvc dies. It passes in the real env.

🤖 Generated with Claude Code

`dt get` had no distributed mode. `dt push --workers` had solved multi-node
transfer already, and the machinery was never push-specific, so this reuses it
rather than rewriting it (#172).

Scoped to `s3://` destinations, which is a decision rather than an unfinished
edge:

  -o s3://...     bound by the network      -> a second node adds bandwidth
  -o some/dir/    bound by the shared FS    -> a second node adds contention

A local destination on Lustre is limited by the filesystem, not the node, so
`--workers` is a hard usage error there pointing at `-j`. On NCI there is a
second reason to like the S3 case: `copyq` is the queue with outbound network
access, and it is already hpc.DEFAULT_TRANSFER_QUEUE.

Reuse, per the issue:

- `hpc.partition_by_size()` is the LPT bin-packer lifted out of
  `push.partition_manifest`, which was pure bin-packing with nothing push-
  specific in it. What stays behind in push is the part that genuinely is: a
  file is named by its hash, so its size has to be read from the local cache.
  The docstring says where the boundary is and why LPT belongs to qxub workers
  and not to a ThreadPoolExecutor -- separate processes cannot steal work, a
  shared queue self-corrects, and freezing an assignment the queue would fix is
  strictly worse.
- `hpc.save_manifest()` takes a `metadata` dict. Everything a worker needs
  beyond its own partition travels that way, because it is a fresh process on
  another machine: push needs a remote name, get also needs a destination,
  credentials config and the flags the run was invoked with.
- `submit_workers` / `monitor_jobs` / `load_worker_partition` unchanged.

The design questions #172 raised, settled:

*Reachability.* Resolve once on the submitting node and pass paths. The
submitter has to resolve every row anyway -- LPT needs byte sizes and
`list_source_files` is where sizes come from -- so having paid for that it ships
the resolved md5, size and destination key per file. **No worker ever runs
`dvc list`**, which is what keeps N compute nodes off the shared clone's SQLite
state db; `_run_dvc_list`'s retry loop exists because that lock is real. A
worker opens the clone once, read-only, for one thing: the source remote's
filesystem and credentials out of its `.dvc/config`.

Link types and the destination identity are likewise resolved on the submitter
and passed through, not re-derived per worker. `resolve_link_types()` reads
`dvc config` from the cwd and `preflight()` resolves a credential chain; both
could answer differently on a compute node, and a transfer that silently
changes its mind halfway is worse than one that fails.

*Partition unit.* Per file, not per row -- the issue framed this as a trade
against per-row reporting, but it is not one. Each task carries its row, so
reporting stays per-row and in CSV order even when a sample directory spans
nodes, and per-file balances properly in the two cases per-row cannot: a single
row of thousands of files, and rows differing by an order of magnitude, which
fastq samples reliably do.

*Reporting.* Workers write `result_<n>.json`; the submitter collates by row tag.
A row whose files do not all come back is failed even when every result that
did arrive succeeded -- a job killed on walltime would otherwise read as a clean
run of a short row.

Each worker preflights for itself, so `--dest-account-id` is enforced on the
node that actually writes. The manifest carries the profile *name*, endpoint and
region; no key material.

Also, from the issue's note about the ThreadPoolExecutor idiom being copy-pasted
five times: this adds a sixth use, so per "factor it then" the two in `get.py`
are now behind `_place_all` and `_upload_all`, and `_resolve_rows` is shared by
the local and S3 row-resolve paths. That makes `upload_to_s3` resolve rows
concurrently, which it did not before. A cross-module `utils.run_parallel()` is
still not worth it.

Not verified here, and it needs a real submission: whether more nodes actually
move more bytes. #172 asks for that measurement and it remains open -- if the
*source* remote is on /g/data the read side is a shared filesystem again and the
same contention argument applies to it. docs/get.md says so and says to time a
single-node `-j 8` run first.

Tested: 30 new unit tests. The central one replaces `hpc.submit_workers` with a
stub that runs the worker in-process, so the manifest is exercised as a real
interface rather than asserted against its spelling.

Co-Authored-By: Claude <noreply@anthropic.com>
@johnyaku
johnyaku merged commit 3f87c85 into main Aug 20, 2026
1 check passed
@johnyaku
johnyaku deleted the feat/get-workers branch August 20, 2026 14:16
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.

dt get: no distributed mode — reuse hpc.submit_workers + partition_manifest for --workers

1 participant