unique on dask fails to gather results properly
Nobody has claimed this yet.
Assessment
- Difficulty
- 3/5
- Estimated time
- 1-2 days
- Newbie friendliness
- 35/100
- Issue type
- Bug
- Clarity
- Mostly clear
- Activity status
- Stale
- Tech stack
- python
- Domain
- stream-processing
Research direction
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.
Written by the indexing model from the issue text.
Description
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.
- Dominant language
- Python
- Stars
- 1.3k
- Forks
- 149
- PR merge metrics
- No merged PRs in 30d
Contributor guide
First steps
- Read the whole issue, then the project's contributing guide.
- Comment on the issue to say you are picking it up — it saves two people doing the same work.
- Fork the repository and make your change on a branch.
- Open a pull request that references the issue number.
More from python-streamz/streamz
-
Difficulty 3/5 1-2 days Newbie friendliness 35/100
python-streamz/streamz#481 · 4 comments ·
-
Combining the streamz.Stream.filenames() and streamz.Stream.from_textfile() using dask scatter? Open
Difficulty 5/5 Over a week Newbie friendliness 25/100
python-streamz/streamz#480 · 2 comments ·
-
Difficulty 5/5 Over a week Newbie friendliness 20/100
python-streamz/streamz#479 · 1 comment ·
-
Difficulty 5/5 Over a week Newbie friendliness 20/100
python-streamz/streamz#478 · 6 comments ·
-
Difficulty 4/5 3-5 days Newbie friendliness 10/100
python-streamz/streamz#476 · 17 comments · 2 reactions ·
All issues in python-streamz/streamz
Similar issues
-
documentation help wanted
Difficulty 2/5 1-3 hours Newbie friendliness 90/100
-
Difficulty 2/5 1-3 hours Newbie friendliness 90/100
simonw/sqlite-utils#872 ·
-
Difficulty 2/5 1-3 hours Newbie friendliness 88/100
-
Difficulty 2/5 1-3 hours Newbie friendliness 82/100
-
Difficulty 2/5 1-3 hours Newbie friendliness 78/100