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

unique on dask fails to gather results properly

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

还没有人认领这个 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 模板
  • 阅读贡献指南

从这里开始

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

python-streamz/streamz 的其他 Issue

查看 python-streamz/streamz 的全部 Issue

相似的 Issue

更多 Python Issue

把新 issue 发到你的邮箱

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