vishalmore90 opened a new pull request, #40283: URL: https://github.com/apache/beam/pull/40283
### Context Fixes #40282. The experimental Python Dask Runner processes side inputs using `DaskBagWindowedIterator`. Previously, this iterator materialized the entire side input dataset directly into client memory by calling `list(self.bag)`. This triggered a full `.compute()` on the Dask Bag, overriding Dask's distributed nature, causing Out-Of-Memory (OOM) errors on large datasets, and preventing pipelines from scaling. ### Changes * **`sdks/python/apache_beam/runners/dask/transform_evaluator.py`**: Refactored the `__iter__` method in `DaskBagWindowedIterator`. * Removed the blocking `list(self.bag)` call. * Utilized `self.bag.to_delayed()` to fetch delayed partition objects and incrementally `.compute()` each partition chunk. This allows the Garbage Collector to clean up processed chunks, keeping the memory footprint constrained to a single partition at a time. ### Verification * Manually validated that `partition.compute()` correctly streams sub-results of the `Bag` iteratively. * Ensured windowing assignments and tagged values remain unaffected by the chunking process. * Ran Dask runner specific tests in the Python SDK. ### PR Checklist - [x] I have read the [CONTRIBUTING.md](https://github.com/apache/beam/blob/master/CONTRIBUTING.md) and [Code of Conduct](https://github.com/apache/beam/blob/master/CODE_OF_CONDUCT.md). - [x] I have analyzed the root cause and implemented a minimal, safe fix. - [x] The fix addresses the exact issue without regressions or style violations. - [x] No additional tests were strictly necessary as this is an under-the-hood memory optimization, and existing Dask test suites cover the correctness. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
