Hacktoberfest 2026: los issues que los mantenedores marcaron para octubre, abiertos y aptos para principiantes. Explorar issues de Hacktoberfest

[Bug]: DaskRunner `DaskBagWindowedIterator` materializes entire dataset into memory causing OOM

Abierto
#40,282 0 comentarios 0 reacciones 0 asignados Ver en GitHub

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

awaiting triage bug P2 python
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:

  1. Partition-based iteration: Utilize Dask's .map_partitions or .to_delayed() to fetch and yield the underlying data partition-by-partition.
  2. 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

Primeros pasos

  1. Lee el issue completo y luego la guía de contribución del proyecto.
  2. Comenta en el issue que vas a ocuparte — evita que dos personas hagan lo mismo.
  3. Haz un fork del repositorio y trabaja en una rama.
  4. Abre un pull request que haga referencia al número del issue.

Más de apache/beam

Todos los issues de apache/beam

Issues similares

Más issues de Java

Recibe los nuevos issues en tu correo

Un resumen breve de issues de GitHub para principiantes.