【问题标题】:Asyncio print status of coroutines progress协程进度的异步打印状态
【发布时间】:2015-07-10 14:06:11
【问题描述】:

我有一堆协程在做一些工作

@asyncio.coroutine
def do_work():
    global COUNTER
    result = ...
    if result.status == 'OK':
        COUNTER += 1

还有一个

COUNTER = 0
@asyncio.coroutine
def display_status():
    while True:
        print(COUNTER)
        yield from asyncio.sleep(1)

必须显示有多少协程完成了他们的工作。如何正确执行此任务?以下解决方案不起作用

@asyncio.coroutine
def spawn_jobs():
    coros = []
    for i in range(10):
        coros.append(asyncio.Task(do_work()))
    yield from asyncio.gather(*coros)

if __name__ == '__main__':
    loop = asyncio.get_event_loop()
    loop.create_task(display_status())
    loop.run_until_complete(spawn_jobs())
    loop.close()

我希望无论 do_work() 协程做什么,计数器都会每秒打印到控制台。但我只有两个输出:0 和几秒钟后重复 10。

【问题讨论】:

  • 一种可能的选择是将 print(COUNTER) 添加到 do_work() 函数中,但我希望每秒打印一次 COUNTER 而不是在完成工作后
  • “不起作用”不是很具体。你期望会发生什么?会发生什么?
  • 不要在 cmets 中添加其他信息,edit 你的问题。
  • 不相关:使用asyncio.create_task() 而不是asyncio.Task()。此外,在这种情况下,您可能都不需要:yield from asyncio.wait([do_work() for _ in range(10)]) 按原样工作。

标签: python python-3.x asynchronous coroutine


【解决方案1】:

但我只有两个输出:0 和几秒钟后重复 10。

我无法复制它。如果我使用:

import asyncio
import random

@asyncio.coroutine
def do_work():
    global COUNTER
    yield from asyncio.sleep(random.randint(1, 5))
    COUNTER += 1

我得到这样的输出:

0
0
4
6
8
Task was destroyed but it is pending!
task: <Task pending coro=<display_status() running at most_wanted.py:16> wait_for=<Future pending cb=[Task._wakeup()] created at ../Lib/asyncio/tasks.py:490> created at most_wanted.py:27>

display_status() 中的无限循环导致最后的警告。避免警告;批处理中的所有任务都完成后退出循环:

#!/usr/bin/env python3
import asyncio
import random
from contextlib import closing
from itertools import cycle

class Batch:
    def __init__(self, n):
        self.total = n
        self.done = 0

    async def run(self):
        await asyncio.wait([batch.do_work() for _ in range(batch.total)])

    def running(self):
        return self.done < self.total

    async def do_work(self):
        await asyncio.sleep(random.randint(1, 5)) # do some work here
        self.done += 1

    async def display_status(self):
        while self.running():
            await asyncio.sleep(1)
            print('\rdone:', self.done)

    async def display_spinner(self, char=cycle('/|\-')):
        while self.running():
            print('\r' + next(char), flush=True, end='')
            await asyncio.sleep(.3)

with closing(asyncio.get_event_loop()) as loop:
    batch = Batch(10)
    loop.run_until_complete(asyncio.wait([
        batch.run(), batch.display_status(), batch.display_spinner()]))

输出

done: 0
done: 2
done: 3
done: 4
done: 10

【讨论】:

  • 问题在于在do_work() 中使用同步代码。我使用同步调用r = requests.get('http://httpbin.org/get'),所以它们会阻止我的显示状态协程。好像我需要切换到 aiohttp 来发出请求。
  • 当我插入 r = yield from aiohttp.request('POST', 'http://httpbin.org/post', data=data) 而不是 sleep() 时,我只能执行少于 512 个批处理作业。否则我有一个错误:RuntimeError: Event loop is closed。是 aiohttp 限制还是别的什么?
  • @MostWanted:你应该问一个单独的问题。 Provide a minimal but complete code example,启用asyncio debug mode,包括完整的回溯。我猜,你需要result = yield from r.json() 之类的东西来阅读回复。
【解决方案2】:

使用threading模块的解决方案

SHOW_STATUS = True
def display_status_sync():
    while SHOW_STATUS:
        print(S_REQ)
        time.sleep(1)

if __name__ == '__main__':
    new_thread = threading.Thread(target=display_status_sync)
    new_thread.start()
    loop.run_until_complete(spawn_jobs())
    SHOW_STATS = False
    new_thread.join()
    loop.close()

但我想使用 asyncio 协程实现类似的功能。有可能吗?

【讨论】:

  • 我也可以使用loop.run_in_executor(None, display_status_sync) 而不是创建线程。这将产生相同的效果。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2023-01-03
  • 1970-01-01
  • 1970-01-01
  • 2016-11-22
  • 1970-01-01
  • 2014-04-29
相关资源
最近更新 更多