Simple example of stream with Asyncio operations
还没有人认领这个 Issue。
评估
- 难度
- 4/5
- 预计耗时
- 3-5 天
- 新手友好度
- 25/100
- Issue 类型
- 缺陷
- 描述清晰度
- 需要澄清
- 活跃度
- 停滞
- 技术栈
- python
调研方向
Start by reproducing the Python example from the issue with its aiohttp, asyncio, and Stream(asynchronous=True) pipeline. Read the streamz asyncio-related implementation and existing tests, if present, to determine whether async map operations are supported. Done means the reported pipeline runs correctly or the supported limitation is documented.
由索引模型根据 Issue 内容生成。
描述
Hi,
I've been following the docs and reading through the tests, and I cannot get streamz working with Asyncio :/
Here's a very minimal example of a stream comprising of two async operation and one sync :
- we retrieve content via an aiohttp call
- return content length
- simulate DB write with an Asyncio sleep
- sink to stdout
import asyncio
import aiohttp
from streamz import Stream
async def fetch(url):
print("fetching url {}", url)
async with aiohttp.ClientSession() as session:
async with session.get(url) as resp:
print(resp.status)
body = await resp.text()
print("Finished w/ url {}", url)
return body
def count(x):
print("I", x)
return len(x)
async def write(x):
await asyncio.sleep(0.2)
print("O", x)
return x
async def f():
print("Starting stream")
source = Stream(asynchronous=True)
source.map(fetch).map(count).map(write).sink(print)
urls = [
'https://httpstatus.io/?i=1',
'https://httpstatus.io/?i=2',
'https://httpstatus.io/?i=3',
'https://httpstatus.io/?i=4',
'https://httpstatus.io/?i=5',
'https://httpstatus.io/?i=6',
]
for u in urls:
await source.emit(u)
if __name__ == '__main__':
asyncio.run(f())
I've tried a lot of combinations using tornado event loop etc. but didn't manage to get anything working.
Is this supposed to be possible or is the Asyncio support still behind?
Am I missing something obvious?
Thanks for the help
- 主要语言
- 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
-
#bug
难度 1/5 1 小时以内 新手友好度 92/100
apache/superset#44923 · 1 条评论 ·
维护者通常 2 天内回复
-
难度 2/5 1-3 小时 新手友好度 84/100
lawndoc/stack-back#123 ·
-
Add: Entune未关闭
难度 2/5 1-3 小时 新手友好度 76/100
AbdelStark/awesome-typesafe-jev#187 ·
维护者通常 1 天内回复
-
bug good first issue
难度 2/5 1-3 小时 新手友好度 88/100
repowise-dev/repowise#2966 · 1 条评论 ·
维护者通常 1 天内回复
-
难度 2/5 1-3 小时 新手友好度 84/100
维护者通常 2 天内回复