Repository navigation
Conversation
|
Assigning reviewers: R: @claudevdm for label python. Note: If you would like to opt out of this review, comment Available commands:
The PR bot will only process comments in the main thread (not review comments). |
|
R: @tvalentyn Hi @tvalentyn, I reviewed your feedback on #40283, particularly the concerns around nested Dask computations, potential worker deadlocks, and changing the default behavior. I've updated this PR to make partition-at-a-time evaluation opt-in, preserving the existing default behavior. Both paths now use Would you be willing to review the updated approach when you have time? I'd appreciate your feedback. Thanks! |
|
Stopping reviewer notifications for this pull request: review requested by someone other than the bot, ceding control. If you'd like to restart, comment |
Summary
Add an opt-in
--dask_lazy_side_inputsmode that consumes Dask Bag side inputs one partition at a time. Keep whole-bag evaluation as the default for compatibility. When either path runs inside a Dask worker, use Dask'sworker_client()while waiting for nested computation so the worker task releases its scheduling slot.Problem and root cause
DaskBagWindowedIteratororiginally evaluatedlist(self.bag)before yielding a value. This can materialize a large side input in one worker and exhaust its memory. Simply replacing it with sequentialpartition.compute()calls also introduces one scheduler round trip per partition and can block a worker task while its nested Dask tasks wait for a free slot. A one-worker, one-thread cluster is especially vulnerable to scheduler starvation.Design
worker_client()temporarily releases the task slot during computation. Its context ends before values are yielded, including when a consumer stops early.ValueErrorfromget_worker()to identify a non-worker context.lazy_side_inputsfrom its options before constructingdask.distributed.Client.This follows Dask's tasks-launching-tasks guidance. It also responds to the concerns raised on the earlier, closed PR #40283, including default-path safety, serial partition evaluation, and exception handling.
Testing
LocalClusterandClient: 4 passed. End-to-end Beam pipelines exerciseAsIter,AsList, andAsSingletonwith one worker/one thread and two workers. A worker-side later-partition error also propagates. Each distributed case runs in a separate process with a 45-second watchdog.dask_runner_test.py: 36 passed, 1 existing expected failure, 5 subtests passed. The run reported 22 warnings, including a missingpytest-timeoutplugin configuration and Dask local dashboard port warnings.yapf --diffon the three changed files,python -m compileall -qon those files, andgit diff --check: passed.The tests ran with Beam's Python SDK and Dask 2026.8.0 installed in an isolated virtual environment. The distributed tests use Dask's distributed scheduler with in-process workers; they do not cover a remote multi-host cluster.
Compatibility and limitations
The default still materializes the entire side-input bag, so the issue's memory limitation remains for users who do not opt in. In opt-in mode, peak iterator materialization is bounded by the largest computed partition, but
AsListmust still materialize its complete view. Partitions are computed serially and may incur repeated graph work or scheduler round trips when views are consumed again or by multiple worker tasks. A shared or storage-backed side-input design would require separate design work.Related to #40282