Hacktoberfest 2026:维护者为十月标记出来的 issue,仍然开放、适合新手。 浏览 Hacktoberfest issue

Flatten doesn't work with DaskStream

未关闭
#213 5 条评论 0 个 reaction 已指派 0 人 在 GitHub 查看

还没有人认领这个 Issue。

评估

难度
3/5
预计耗时
1-2 天
新手友好度
35/100
Issue 类型
缺陷
描述清晰度
基本清楚
活跃度
停滞
技术栈
python

调研方向

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.

由索引模型根据 Issue 内容生成。

描述

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.

主要语言
Python
星标
1.3k
派生
151
平均合并
17 小时 39 分钟
30 天内合并 PR
1

环境准备

  • 提供 Dockerfile 或 Docker Compose 文件
  • 没有 Pull Request 模板
  • 阅读贡献指南

从这里开始

  1. 先读完整个 Issue,再读项目的贡献指南。
  2. 在 Issue 下留言说明你要接手 —— 这能避免两个人做同样的事。
  3. Fork 仓库,在一个分支上完成修改。
  4. 提交 Pull Request,并在描述里引用这个 Issue 编号。

python-streamz/streamz 的其他 Issue

查看 python-streamz/streamz 的全部 Issue

相似的 Issue

更多 Python Issue

把新 issue 发到你的邮箱

精选适合新手参与的 GitHub issue 摘要。