【问题标题】:Fire coroutine from instide a for loop从 for 循环内部触发协程
【发布时间】:2019-10-15 11:18:18
【问题描述】:

我正在尝试从循环中触发协程。这是我想要实现的一个简单示例:

import time
import random
import asyncio


def listen():
    while True:
       yield random.random()
       time.sleep(3)


async def dosomething(data: float):
    print("Working on data", data)
    asyncio.sleep(2)
    print("Processed data!")


async def main():
    for pos in listen():
        asyncio.create_task(dosomething(pos))


if __name__ == '__main__':
    loop = asyncio.get_event_loop()
    loop.run_until_complete(main())

不幸的是,这不起作用,我的 dosomething 协程永远不会执行......我做错了什么?

【问题讨论】:

  • Async 还会抱怨main 在释放控制权之前花费了太多时间。这只是一个警告,但值得在您的实际应用中考虑。

标签: python python-asyncio coroutine


【解决方案1】:

asyncio.create_task 函数旨在调度任务执行,它应该等待直到它完成。

此外,您的代码中的asyncio.sleep(2) 也应该等待,否则它会引发错误/警告。

正确的方法:

import time
import random
import asyncio


def listen():
    while True:
       yield random.random()
       time.sleep(3)


async def dosomething(data: float):
    print("Working on data", data)
    await asyncio.sleep(2)
    print("Processed data!")


async def main():
    for pos in listen():
        await asyncio.create_task(dosomething(pos))


if __name__ == '__main__':
    loop = asyncio.get_event_loop()
    loop.run_until_complete(main())

示例输出:

Working on data 0.9645515392725723
Processed data!
Working on data 0.9249656672476657
Processed data!
Working on data 0.13635467058997397
Processed data!
Working on data 0.03941252405458562
Processed data!
Working on data 0.6299882183389822
Processed data!
Working on data 0.9143748948769984
Processed data!
...

【讨论】:

  • 感谢@RomanPerekhrest - 我正在尝试像在 javascript 中一样使用 async/await >.
  • 当您对任务对象除了等待它之外不做任何事情时,创建任务是不必要的开销。你可以简单地await dosomething(pos) 代替。
  • @NickMartin 请注意,此代码将按顺序执行协程。如果您打算并行执行它们,您应该调用类似await asyncio.gather(*[dosomething(pos) for pos in listen()])
  • 另外,listen() 应该是一个async def 并使用await asyncio.sleep(),否则它对time.sleep() 的调用将继续阻塞整个事件循环。然后main() 应该使用async for 对其进行迭代。
【解决方案2】:

我想指出,在玩过之后,我最终使用了生产者、消费者架构来实现我想要的。我很感激我没有在原始问题中明确说明我的确切用例。但这是我最终实现的简化的 sn-p:

import asyncio
import random
from datetime import datetime
from pydantic import BaseModel


class Measurement(BaseModel):
    data: float
    time: datetime


async def measure(queue: asyncio.Queue):
    while True:
        # Replicate blocking call to recieve data
        await asyncio.sleep(1)
        print("Measurement complete!")
        for i in range(3):
            data = Measurement(
                data=random.random(),
                time=datetime.utcnow()
            )
            await queue.put(data)

    await queue.put(None)


async def process(queue: asyncio.Queue):
    while True:
        data = await queue.get()
        print(f"Got measurement! {data}")
        # Replicate pause for http request
        await asyncio.sleep(0.3)
        print("Sent data to server")


loop = asyncio.get_event_loop()
queue = asyncio.Queue(loop=loop)
meansurement = measure(queue)
processor = process(queue)
loop.run_until_complete(asyncio.gather(processor, meansurement))
loop.close()

我应该在这里指出(我不太明白的一点),你所做的任何阻塞调用都必须是await-ed。否则,你可能会发现消费者永远不会执行。

【讨论】:

    猜你喜欢
    • 2015-03-21
    • 1970-01-01
    • 1970-01-01
    • 2017-09-20
    • 2020-04-04
    • 1970-01-01
    • 2014-03-05
    • 2013-06-07
    • 2011-03-23
    相关资源
    最近更新 更多