unique on dask fails to gather results properly
还没有人认领这个 Issue。
评估
- 难度
- 3/5
- 预计耗时
- 1-2 天
- 新手友好度
- 35/100
- Issue 类型
- 缺陷
- 描述清晰度
- 基本清楚
- 活跃度
- 停滞
- 技术栈
- python
调研方向
Run the supplied DaskStream reproducer, comparing scatter().unique().buffer(10).gather() with the working Stream and scatter().buffer(10).gather() paths. Trace the unique and gather entry points and verify that the gathered sink_to_list output contains plain values rather than futures.
由索引模型根据 Issue 内容生成。
描述
It seems that when gathering after a unique on a DaskStream the results fail to be returned to their non future state. However, this seems limited to only unique on dask, running either with unique in a standard Stream or without in a DaskStream seems to do the trick.
Working example:
In [1]: from dask.distributed import Client
...: client = Client()
...:
In [2]: from streamz import Stream
In [3]: s = Stream()
In [4]: b = s.unique()
In [5]: c = s.scatter().unique().buffer(10).gather()
In [6]: l = b.sink_to_list()
In [8]: ll = c.sink_to_list()
In [9]: for i in range(10):
...: s.emit(i)
In [10]: l
Out[10]: [0, 1, 2, 3, 4, 5, 6, 7, 8, 9]
In [11]: ll
Out[11]:
[<Future: status: finished, type: int, key: int-5c8a950061aa331153f4a172bbcbfd1b>,
<Future: status: finished, type: int, key: int-c0a8a20f903a4915b94db8de3ea63195>,
<Future: status: finished, type: int, key: int-58e78e1b34eb49a68c65b54815d1b158>,
<Future: status: finished, type: int, key: int-d3395e15f605bc35ab1bac6341a285e2>,
<Future: status: finished, type: int, key: int-5cd9541ea58b401f115b751e79eabbff>,
<Future: status: finished, type: int, key: int-ce9a05dd6ec76c6a6d171b0c055f3127>,
<Future: status: finished, type: int, key: int-7ec5d3339274cee5cb507a4e4d28e791>,
<Future: status: finished, type: int, key: int-06e5a71c9839bd98760be56f629b24cc>,
<Future: status: finished, type: int, key: int-ea1fa36eb048f89cc9b6b045a2a731d2>,
<Future: status: finished, type: int, key: int-c56e7bae3484c9b6750417fbf89d6509>]
In [12]: d = s.scatter().buffer(10).gather()
In [13]: lll = d.sink_to_list()
In [14]: for i in range(10):
...: s.emit(i)
...:
In [15]: lll
Out[15]: [0, 1, 2, 3, 4, 5, 6, 7, 8, 9]
This was produced with Python 3.5.5 | packaged by conda-forge | (default, Feb 13 2018, 05:02:37) and latest master Streamz.
- 主要语言
- Python
- 星标
- 1.3k
- 派生
- 151
- 平均合并
- 17 小时 39 分钟
- 30 天内合并 PR
- 1
环境准备
- 提供 Dockerfile 或 Docker Compose 文件
- 没有 Pull Request 模板
- 阅读贡献指南
从这里开始
- 先读完整个 Issue,再读项目的贡献指南。
- 在 Issue 下留言说明你要接手 —— 这能避免两个人做同样的事。
- Fork 仓库,在一个分支上完成修改。
- 提交 Pull Request,并在描述里引用这个 Issue 编号。
python-streamz/streamz 的其他 Issue
-
难度 3/5 1-2 天 新手友好度 35/100
python-streamz/streamz#481 · 4 条评论 ·
-
难度 5/5 一周以上 新手友好度 25/100
python-streamz/streamz#480 · 2 条评论 ·
-
难度 5/5 一周以上 新手友好度 20/100
python-streamz/streamz#479 · 1 条评论 ·
-
难度 5/5 一周以上 新手友好度 20/100
python-streamz/streamz#478 · 6 条评论 ·
-
难度 4/5 3-5 天 新手友好度 10/100
python-streamz/streamz#476 · 17 条评论 · 2 个 reaction ·
查看 python-streamz/streamz 的全部 Issue
相似的 Issue
-
tool-calling
难度 2/5 1-3 小时 新手友好度 88/100
vllm-project/vllm#59838 ·
维护者通常 1 天内回复
-
难度 1/5 1 小时以内 新手友好度 92/100
raullenchai/Rapid-MLX#4042 ·
维护者通常 1 天内回复
-
documentation
难度 1/5 1 小时以内 新手友好度 92/100
transitmatters/mbta-slow-zone-bot#70 ·
维护者通常 1 天内回复
-
bug
难度 2/5 1-3 小时 新手友好度 78/100
litestar-org/advanced-alchemy#811 ·
维护者通常 1 天内回复
-
难度 2/5 1-3 小时 新手友好度 78/100
维护者通常 1 天内回复