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

Flatten doesn't work with DaskStream

Abierto
#213 5 comentarios 0 reacciones 0 asignados Ver en GitHub

Nadie ha tomado este issue todavía.

Evaluación

Dificultad
3/5
Tiempo estimado
1-2 días
Aptitud para principiantes
35/100
Tipo de issue
Error
Claridad
Bastante claro
Estado de actividad
Estancado
Stack tecnológico
python

Línea de trabajo

Reproduce the provided script, then inspect streamz/core.py around the flatten update at line 1020 and streamz/dask.py around line 58, where the traceback shows a Future being passed onward. Done means the scatter, partition, map, flatten, buffer, gather, and sink pipeline no longer raises the reported TypeError.

Escrito por el modelo de indexación a partir del texto del issue.

Descripción

This code:

from __future__ import division, print_function

from time import sleep
from streamz import Stream
from dask.distributed import Client

client = Client()

def callback(datas):
    print(':',datas)
    return datas

source = Stream().scatter()
stream = source.partition(5)
stream = stream.map(callback)
stream = stream.flatten()

stream.buffer(15).gather().sink(print)

for i in range(30):
    source.emit(i)

returns:

Traceback (most recent call last):
  File "test.py", line 21, in <module>
: (0, 1, 2, 3, 4)
    source.emit(i)
  File "/home/klinger/years/2014/projects/devel/streamz/streamz/core.py", line 332, in emit
    sync(self.loop, _)
  File "/home/klinger/years/2014/projects/devel/streamz/streamz/core.py", line 1277, in sync
    six.reraise(*error[0])
  File "/home/klinger/years/2014/projects/devel/streamz/streamz/core.py", line 1262, in f
    result[0] = yield future
  File "/home/klinger/rossum/local/lib/python2.7/site-packages/tornado/gen.py", line 1099, in run
    value = future.result()
  File "/home/klinger/rossum/local/lib/python2.7/site-packages/tornado/concurrent.py", line 260, in result
    raise_exc_info(self._exc_info)
  File "/home/klinger/rossum/local/lib/python2.7/site-packages/tornado/gen.py", line 315, in wrapper
    yielded = next(result)
  File "/home/klinger/years/2014/projects/devel/streamz/streamz/core.py", line 327, in _
    result = yield self._emit(x)
  File "/home/klinger/years/2014/projects/devel/streamz/streamz/core.py", line 298, in _emit
    r = downstream.update(x, who=self)
  File "/home/klinger/years/2014/projects/devel/streamz/streamz/core.py", line 734, in update
    return self._emit(tuple(result))
  File "/home/klinger/years/2014/projects/devel/streamz/streamz/core.py", line 298, in _emit
    r = downstream.update(x, who=self)
  File "/home/klinger/years/2014/projects/devel/streamz/streamz/dask.py", line 58, in update
    return self._emit(result)
  File "/home/klinger/years/2014/projects/devel/streamz/streamz/core.py", line 298, in _emit
    r = downstream.update(x, who=self)
  File "/home/klinger/years/2014/projects/devel/streamz/streamz/core.py", line 1020, in update
    for item in x:
TypeError: 'Future' object is not iterable

Without dask interface (removing scatter and gather), everything works well.

Lenguaje dominante
Python
Estrellas
1.3k
Forks
149
Merge medio
17 h 39 min
PR fusionados (30 d)
1

Guía de contribución

Abrir la guía de contribución

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 python-streamz/streamz

Todos los issues de python-streamz/streamz

Issues similares

Más issues de Python

Recibe los nuevos issues en tu correo

Un resumen breve de issues de GitHub para principiantes.