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]

Reply via email to