vishalmore90 opened a new issue, #40243:
URL: https://github.com/apache/beam/issues/40243

   ### What happened?
   
   
   In the Spark Structured Streaming runner, `ParDoTranslatorBatch.java` 
explicitly calls `.persist()` on intermediate datasets when a `ParDo` transform 
has multiple outputs (side outputs). This is correctly done to avoid 
re-evaluating the parent `mapPartitions` transformation multiple times when 
dividing the dataset by tuple tags. However, the runner never calls 
`.unpersist()` to free these datasets from Spark's storage memory. 
   
   This creates a significant memory leak. The materialized datasets will 
accumulate in Spark's block manager indefinitely until the Spark application 
terminates (or they are forcefully evicted by Spark's LRU cache, which degrades 
performance).
   
   **Code Pointers & Steps to Reproduce:**
   1. Navigate to 
`runners/spark/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/ParDoTranslatorBatch.java`
 (around line 138-146).
   2. The code explicitly acknowledges this technical debt with a comment: `// 
FIXME What's the strategy to unpersist Datasets / RDDs?`
   3. Immediately following the comment, `allTagsDS.persist(storageLevel);` is 
called to cache the dataset, but no corresponding `unpersist()` logic exists 
anywhere in the translation lifecycle.
   4. **To reproduce:** Run a large Beam batch pipeline with the Spark 
Structured Streaming runner that contains several multi-output `ParDo` 
transforms. Monitor Spark's Storage memory tab in the Spark UI. The cached 
datasets will continue to grow without ever being released.
   
   **Impact:**
   - Causes Out-Of-Memory (OOM) errors and uncontrolled disk spilling in Spark 
clusters due to exhaustion of Spark's storage memory.
   - Major performance degradation for large batch pipelines relying on the 
Spark Structured Streaming runner.
   
   **Proposed Solution:**
   - Implement a cleanup strategy/lifecycle hook for unpersisting intermediate 
datasets when they are no longer needed by downstream consumers in the Spark 
Structured Streaming execution plan. 
   - Alternatively, attach a Spark `SparkListener` or use a reference counting 
mechanism within the translation context to explicitly unpersist these datasets 
once all downstream stages consuming the side outputs have completed.
   
   
   
   
   
   
   ### Issue Priority
   
   Priority: 2 (default / most bugs should be filed as P2)
   
   ### Issue Components
   
   - [ ] Component: Python SDK
   - [ ] Component: Java SDK
   - [ ] Component: Go SDK
   - [ ] Component: Typescript SDK
   - [ ] Component: IO connector
   - [ ] Component: Beam YAML
   - [ ] Component: Beam examples
   - [ ] Component: Beam playground
   - [ ] Component: Beam katas
   - [ ] Component: Website
   - [ ] Component: Infrastructure
   - [x] Component: Spark Runner
   - [ ] Component: Flink Runner
   - [ ] Component: Prism Runner
   - [ ] Component: Twister2 Runner
   - [ ] Component: Hazelcast Jet Runner
   - [ ] Component: Google Cloud Dataflow Runner


-- 
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