[Bug]: DaskRunner `DaskBagWindowedIterator` materializes entire dataset into memory causing OOM
I maintainer di solito rispondono entro 1 giorno
@Gaurav598 ci sta già lavorando.
Dal 10/10/2026.
Valutazione
- Difficoltà
- 4/5
- Tempo stimato
- 3-5 giorni
- Idoneità per principianti
- 56/100
- Tipo di issue
- Bug
- Chiarezza
- Abbastanza chiara
- Stato di attività
- Attiva
- Stack tecnologico
- python
- Ambito
- data-engineering, distributed-systems
Direzione di ricerca
Start in sdks/python/apache_beam/runners/dask/transform_evaluator.py at DaskBagWindowedIterator.iter and review the FIXME around list(self.bag). Investigate the partition-based and delayed-result approaches described in the issue. Done means side-input results are yielded incrementally without materializing the full Dask Bag in client memory, avoiding the reported OOM behavior.
Scritto dal modello di indicizzazione a partire dal testo della issue.
Descrizione
What happened?
In the experimental Python Dask Runner, the DaskBagWindowedIterator (which handles iterators for apache_beam.transforms.sideinputs.SideInputMap) iterates over a Dask Bag by wrapping it in a Python list().
Calling list(self.bag) implicitly triggers a full compute() on the Dask dataset. This blocking operation materializes the entirety of the side input data into the client's local memory. For large side inputs, this completely bypasses Dask's distributed memory management and results in an Out-Of-Memory (OOM) crash, effectively bottlenecking the scalability of pipelines running on Dask.
The code currently includes an explicit FIXME acknowledging this proof-of-concept behavior, but it remains a silent, critical scalability flaw.
Code Pointers / Steps to Reproduce
The issue is located in sdks/python/apache_beam/runners/dask/transform_evaluator.py within the __iter__ method of the DaskBagWindowedIterator class (lines 91-96):
class DaskBagWindowedIterator:
"""Iterator for `apache_beam.transforms.sideinputs.SideInputMap`"""
bag: db.Bag
window_fn: WindowFn
def __iter__(self):
# FIXME(cisaacstern): list() is likely inefficient, since it presumably
# materializes the full result before iterating over it. doing this for
# now as a proof-of-concept. can we can generate results incrementally?
for result in list(self.bag):
yield get_windowed_value(result, self.window_fn)
Impact
Any Apache Beam pipeline using DaskRunner that relies on substantial side inputs will crash with OOM errors as soon as the side input data surpasses the available local RAM on the node where the iterator is evaluated. This severely limits the DaskRunner's ability to process real-world distributed datasets and creates a harsh scalability ceiling.
Proposed Solution
The evaluation should generate results incrementally rather than performing a monolithic evaluation. Potential approaches:
- Partition-based iteration: Utilize Dask's
.map_partitionsor.to_delayed()to fetch and yield the underlying data partition-by-partition. - Generators: Instead of eager computation via
list(), retrieve delayed results asynchronously and yield them to allow the Python garbage collector to free memory between partition iterations.
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
- Component: Spark Runner
- Component: Flink Runner
- Component: Prism Runner
- Component: Twister2 Runner
- Component: Hazelcast Jet Runner
- Component: Google Cloud Dataflow Runner
- Lingua principale
- Java
- Stelle
- 8.7k
- Fork
- 4.7k
- Merge medio
- 2g 7h
- PR unite (30g)
- 242
Preparare l'ambiente
- Nessun Dockerfile né file Docker Compose
- Ha un modello di pull request
- Leggi la guida per i contributori
Come iniziare
- Leggi tutta la issue e poi la guida ai contributi del progetto.
- Commenta sulla issue per dire che te ne occupi tu — evita che due persone facciano lo stesso lavoro.
- Fai un fork del repository e lavora su un branch.
- Apri una pull request che faccia riferimento al numero della issue.
Altre issue di apache/beam
-
[Bug]: Row.toString throws for an ITERABLE field that is not backed by a ListForse già presa @PDGGK l’ha presa 58 giorni fa. Apertajava P3
Difficoltà 2/5 1-3 ore Idoneità per principianti 76/100
I maintainer di solito rispondono entro 1 giorno
-
Difficoltà 2/5 1-3 ore Idoneità per principianti 68/100
apache/beam#39624 · 2 reazioni ·
I maintainer di solito rispondono entro 1 giorno
-
[Bug]: PubsubIO used in batch incorrect batch cutoff sizeForse già presa @1fanwang l’ha presa 46 giorni fa. Apertabug io P3 pinned pubsub
Difficoltà 2/5 1-3 ore Idoneità per principianti 78/100
apache/beam#28011 · 4 commenti ·
I maintainer di solito rispondono entro 1 giorno
-
build P3 sub-task
Difficoltà 1/5 Meno di un'ora Idoneità per principianti 68/100
I maintainer di solito rispondono entro 1 giorno
-
bug gcp io java P3
Difficoltà 1/5 Meno di un'ora Idoneità per principianti 65/100
I maintainer di solito rispondono entro 1 giorno
Issue simili
-
BoxAttachmentMulti parsing leaks IOException / ArrayIndexOutOfBoundsException on malformed content instead of IllegalArgumentExceptionForse già presa @Kshot3000 l’ha presa oggi. Aperta
Difficoltà 2/5 1-3 ore Idoneità per principianti 74/100
ergoplatform/ergo-appkit#272 ·
-
Difficoltà 2/5 1-3 ore Idoneità per principianti 64/100
-
Difficoltà 2/5 1-3 ore Idoneità per principianti 66/100
-
Difficoltà 2/5 1-3 ore Idoneità per principianti 64/100
utopia-rise/godot-jvm#1004 ·
I maintainer di solito rispondono entro 1 giorno
-
Difficoltà 2/5 1-3 ore Idoneità per principianti 82/100
spring-projects/spring-grpc#442 ·