[Bug]: DaskRunner `DaskBagWindowedIterator` materializes entire dataset into memory causing OOM
Los mantenedores suelen responder en 1 día
Nadie ha tomado este issue todavía.
- #40283 de @vishalmore90 — cerrado sin fusionar
Evaluación
- Dificultad
- 4/5
- Tiempo estimado
- 3-5 días
- Aptitud para principiantes
- 56/100
- Tipo de issue
- Error
- Claridad
- Bastante claro
- Estado de actividad
- Activo
- Stack tecnológico
- python
Línea de trabajo
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.
Escrito por el modelo de indexación a partir del texto del issue.
Descripción
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
- Lenguaje dominante
- Java
- Estrellas
- 8.7k
- Forks
- 4.7k
- Merge medio
- 2 d 8 h
- PR fusionados (30 d)
- 246
Preparar el entorno
- Sin Dockerfile ni archivo de Docker Compose
- Tiene una plantilla de pull request
- Leer la guía de contribución
Primeros pasos
- Lee el issue completo y luego la guía de contribución del proyecto.
- Comenta en el issue que vas a ocuparte — evita que dos personas hagan lo mismo.
- Haz un fork del repositorio y trabaja en una rama.
- Abre un pull request que haga referencia al número del issue.
Más de apache/beam
-
[Bug]: Row.toString throws for an ITERABLE field that is not backed by a ListPosiblemente ocupada @PDGGK la tomó hace 56 días. Abiertojava P3
Dificultad 2/5 1-3 horas Aptitud para principiantes 76/100
Los mantenedores suelen responder en 1 día
-
Dificultad 2/5 1-3 horas Aptitud para principiantes 68/100
apache/beam#39624 · 2 reacciones ·
Los mantenedores suelen responder en 1 día
-
[Failing Test]: JmsIOTest. testCheckpointMark flakyPosiblemente ocupada @mxtymoshyk la tomó hace 13 días. Abiertobug failing test flake P2 pinned tests
Dificultad 2/5 1-3 horas Aptitud para principiantes 68/100
apache/beam#30225 · 2 comentarios ·
Los mantenedores suelen responder en 1 día
-
[Bug]: PubsubIO used in batch incorrect batch cutoff sizePosiblemente ocupada @1fanwang la tomó hace 44 días. Abiertobug io P3 pinned pubsub
Dificultad 2/5 1-3 horas Aptitud para principiantes 78/100
apache/beam#28011 · 4 comentarios ·
Los mantenedores suelen responder en 1 día
-
build P3 sub-task
Dificultad 1/5 Menos de una hora Aptitud para principiantes 68/100
Los mantenedores suelen responder en 1 día
Todos los issues de apache/beam
Issues similares
-
Dificultad 2/5 1-3 horas Aptitud para principiantes 68/100
Netcracker/qubership-integration-platform#1046 ·
Los mantenedores suelen responder en 2 días
-
`check_java_version()` fails when Java path contains spaces (Windows / Git Bash, `C:\Program Files`)Abierto
Dificultad 2/5 1-3 horas Aptitud para principiantes 68/100
-
Fix Math.ceilDiv wrong result for exact positive divisionsPosiblemente ocupada @pamod-madubashana la tomó hoy. Abierto
Dificultad 2/5 1-3 horas Aptitud para principiantes 88/100
scala-native/scala-native#5094 ·
Los mantenedores suelen responder en 1 día
-
[Bug] AI unread message badge counts a batch of new bubbles as one messagePosiblemente ocupada Un pull request vinculado a esta issue está abierto o ya se fusionó. Abierto
Dificultad 2/5 1-3 horas Aptitud para principiantes 74/100
apache/rocketmq-dashboard#5784 ·
Los mantenedores suelen responder en 3 días
-
[i18n] 安装实例完成后的成功提示未正确本地化Abierto
Dificultad 2/5 1-3 horas Aptitud para principiantes 62/100
PCL-Community/PCL-CE#3658 ·
Los mantenedores suelen responder en 1 día