Skip to content

Add opt-in lazy Dask side inputs with worker-safe computation - #40496

Open
Gaurav598 wants to merge 2 commits into
apache:masterfrom
Gaurav598:fix/40282-dask-side-input-partitions
Open

Gaurav598 wants to merge 2 commits into
apache:masterfrom
Gaurav598:fix/40282-dask-side-input-partitions

Conversation

@Gaurav598

@Gaurav598 Gaurav598 commented Oct 10, 2026 •

Copy link
Copy Markdown

Summary

Add an opt-in --dask_lazy_side_inputs mode 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's worker_client() while waiting for nested computation so the worker task releases its scheduling slot.

Problem and root cause

DaskBagWindowedIterator originally evaluated list(self.bag) before yielding a value. This can materialize a large side input in one worker and exhaust its memory. Simply replacing it with sequential partition.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

  • The new flag is off by default. In opt-in mode, the iterator computes one partition when reached and preserves bag order and the existing window conversion.
  • Outside a Dask worker, computation uses the configured Dask scheduler. Inside a worker, worker_client() temporarily releases the task slot during computation. Its context ends before values are yielded, including when a consumer stops early.
  • Computation failures propagate unchanged. The code only catches ValueError from get_worker() to identify a non-worker context.
  • The runner removes lazy_side_inputs from its options before constructing dask.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

  • Focused iterator tests: 5 passed. They cover default whole-bag behavior, partition-at-a-time consumption, empty partitions, repeated iteration, later-partition errors, order, and window conversion.
  • Distributed tests with LocalCluster and Client: 4 passed. End-to-end Beam pipelines exercise AsIter, AsList, and AsSingleton with 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.
  • Full dask_runner_test.py: 36 passed, 1 existing expected failure, 5 subtests passed. The run reported 22 warnings, including a missing pytest-timeout plugin configuration and Dask local dashboard port warnings.
  • yapf --diff on the three changed files, python -m compileall -q on those files, and git 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 AsList must 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

@github-actions

Copy link
Copy Markdown
Contributor

Assigning reviewers:

R: @claudevdm for label python.

Note: If you would like to opt out of this review, comment assign to next reviewer.

Available commands:

  • stop reviewer notifications - opt out of the automated review tooling
  • remind me after tests pass - tag the comment author after tests pass
  • waiting on author - shift the attention set back to the author (any comment or push by the author will return the attention set to the reviewers)

The PR bot will only process comments in the main thread (not review comments).

@Gaurav598 Gaurav598 changed the title Stream Dask side inputs one partition at a time Add opt-in lazy Dask side inputs with worker-safe computation Oct 10, 2026
@Gaurav598

Copy link
Copy Markdown
Author

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 worker_client() for nested computations inside Dask workers, and I've added distributed regression tests covering single-worker/single-thread and multi-worker execution.
The focused and distributed tests pass locally. The PR description also documents the remaining memory and performance trade-offs.

Would you be willing to review the updated approach when you have time? I'd appreciate your feedback. Thanks!

@github-actions

Copy link
Copy Markdown
Contributor

Stopping reviewer notifications for this pull request: review requested by someone other than the bot, ceding control. If you'd like to restart, comment assign set of reviewers

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant