tvalentyn commented on PR #40283: URL: https://github.com/apache/beam/pull/40283#issuecomment-5876770208
looks like code-review bot is not working now; i ran this by an AI offline and got this: ### 👍 Pros * **Resolves the Fatal OOM Bug:** This is the primary victory of the PR. The previous implementation called `list(self.bag)` ([source](https://github.com/apache/beam/pull/40283/changes?content_ref=for+result+in+list+self+bag)), which forced Dask to evaluate the whole graph and pull all results into a single list ([source](https://github.com/apache/beam/issues/40282?content_ref=calling+list+self+bag+implicitly+triggers+a+full+compute+on+the+dask+dataset)), ([source](https://github.com/apache/beam/issues/40282?content_ref=this+blocking+operation+materializes+the+entirety+of+the+side+input+data+into+the+client+s+local+memory)). The PR fixes this by ensuring only a single partition is materialized in Python memory at any given time. * **Leverages Garbage Collection:** By yielding the elements of a single partition and then moving to the next iteration of the loop, Python's garbage collector can clear the previous partition's data out of memory. This keeps the memory high-water mark equal to the size of the single largest partition. * **Low Code Complexity:** The fix is a minimal, three-line change that uses standard Dask primitives without requiring a massive architectural rewrite of how `apache_beam.runners.dask` evaluates side inputs. * **Maintains Determinism:** It preserves the exact sequential order of the original Dask Bag. ### 👎 Cons * **Destroys Parallelism (Serial Execution):** This is the most significant drawback. Because the code loops sequentially (`for partition... partition.compute()`), the Dask cluster will only evaluate **one partition at a time**. If the upstream computation for this side input is heavy, the entire distributed cluster will sit mostly idle while a single worker computes partition 0, then partition 1, and so on. * **"Tasks Launching Tasks" Deadlock Risk:** In Apache Beam, side inputs are typically evaluated inside the `ParDo` worker execution context. Calling a blocking `.compute()` from inside a Dask worker is considered a severe anti-pattern in Dask. If all workers in the cluster are busy executing `ParDo` tasks and they all call `.compute()`, they will wait indefinitely for the scheduler to assign their sub-tasks to a worker thread, causing a cluster-wide deadlock. * **High Scheduler Overhead:** Instead of submitting one large computation graph to the Dask scheduler, this implementation submits $N$ separate computation graphs (where $N$ is the number of partitions). This introduces a round-trip network latency penalty to the scheduler for every single partition. * **The Recomputation Penalty Remains:** The `__iter__` magic method can be invoked multiple times per window or element depending on the Beam transform. Every time `__iter__` is called, this code will re-trigger the entire Dask calculation for that side input from scratch. -- 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]
