【问题标题】:Python asyncio timeout/retry design patternPython asyncio 超时/重试设计模式
【发布时间】:2022-10-08 09:11:55
【问题描述】:

请看下面的代码(为了简单起见,我没有使用 pydantic 来对 corutine、重试、超时进行分组):

import asyncio
import typing as tp
import random

async def my_func(wait_time: int) -> str:
    random_number = random.random()
    random_time = wait_time - random_number if random.random() < 0.5 else wait_time + random_number
    print(f"waiting for {wait_time}{random_time:+} seconds")
    await asyncio.sleep(wait_time)
    return f"waited for {wait_time}{random_time:+} seconds"

async def main() -> None:

    task1 = asyncio.create_task(my_func(wait_time=1), name='task1')
    task2 = asyncio.create_task(my_func(wait_time=2), name='task2')
    task3 = asyncio.create_task(my_func(wait_time=3), name='task3')

    task1_timeout = 1.2
    task2_timeout = 2.2
    task3_timeout = 3.2

    task1_retry = 4
    task2_retry = 3
    task3_retry = 2

    total_timeout = 5

    <what to put here?>

    return task1_result, task2_result, task3_result

asyncio.run(main())

如您所见,我有函数 my_func (在现实生活中我将有多个不同的函数)。 在 main() 中,我定义了 3 个任务。每个任务都有其超时和重试。 例如,task1 超时 2 秒,重试 3 次。

此外,我还有另一个(全局)超时,total_timeout,它表示 main() 必须完成的时间。

例如,如果task1 开始运行并且在 1.2 秒内没有得到结果,我们应该最多重试 4 次,所以在我们根本无法得到结果的情况下,我们仍然低于 timeout_total 的 5秒。

对于task2,在2.2秒内超时,可以重复3次,在4.4秒第二次重复完成后,如果我们再次重试,它将在第5秒被total_timeout截断。

对于task3,如果我们第一次尝试没有完成,我们没有足够的时间进行第二次尝试(total_timeout)。

我想同时执行所有三个任务,尊重他们各自的超时和重试,以及total_timeout。最后,最多 5 秒后,我将得到三个元素的元组,它们将是 str(my_func 的输出)或 None(以防所有重复失败,或任务已被 total_timeout 切断)。 所以输出可以是(str, str, str)(str, None, str)(None, None, None)

有人可以提供一些示例代码来完成我所描述的吗?

【问题讨论】:

  • 你需要像await asyncio.gather(task1, task2, task3) 这样的东西。这将返回三个结果,以便您传入等待对象。但请记住,asyncio 不会同时运行。它允许一个任务在一个或多个其他任务等待 I/O 完成时运行。
  • 收集根本没有超时
  • 你应该使用wait_for 而不是create_task。这几乎是整个timeouts section of the docs
  • 是的,这听起来很容易。你有超时的wait_for(但是一个等待的),你有多个等待超时的等待,你有没有超时的聚集......有很多选择,但我还没有看到有人为什么提供了解决方案我已经描述过了。我认为这是许多人可以从中受益的事情。
  • 您尝试过哪些?他们中的任何一个工作了吗?如果它们不起作用,那么每个版本有什么问题?

标签: python timeout python-asyncio retry-logic


【解决方案1】:

我认为这是一个很好的问题。我提出了这个解决方案,它结合了asyncio.gather()asyncio.wait_for().

这里要求第三个任务等待5秒3.2 秒超时(重试 2 次),并将返回 None,如asyncio.TimeoutError将被提升(并被抓住)。

import asyncio
import random
import sys


total_timeout = float(sys.argv[1]) if len(sys.argv) > 1 else 5.0


async def work_coro(wait_time: int) -> str:
    random_number = random.random()
    random_time = wait_time - random_number if 
        random.random() < 0.5 else wait_time + random_number

    if random_number > 0.7:
        raise RuntimeError('Random sleep time too high')

    print(f"waiting for {wait_time}{random_time:+} seconds")

    await asyncio.sleep(random_time)

    return f"waited for {wait_time}{random_time:+} seconds"


async def coro_trunner(wait_time: int,
                       retry: int,
                       timeout: float) -> str:
    """
    Run work_coro in a controlled timing environment

    :param int wait_time: How long the coroutine will sleep on each run
    :param int retry: Retry count (if the coroutine times out, retry x times)
    :param float timeout: Timeout for the coroutine
    """

    for attempt in range(0, retry):
        try:
            start_time = loop.time()
            print(f'{work_coro}: ({wait_time}, {retry}, {timeout}): '
                  'spawning')

            return await asyncio.wait_for(work_coro(wait_time),
                                          timeout)
        except asyncio.TimeoutError:
            diff_time = loop.time() - start_time
            print(f'{work_coro}: ({wait_time}, {retry}, {timeout}): '
                  f'timeout (diff_time: {diff_time}')
            continue
        except asyncio.CancelledError:
            print(f'{work_coro}: ({wait_time}, {retry}, {timeout}): '
                  'cancelled')
            break
        except Exception as err:
            # Unknown error raised in the worker: give it another chance

            print(f'{work_coro}: ({wait_time}, {retry}, {timeout}): '
                  f'error in worker: {err}')
            continue


async def main() -> list:
    tasks = [
        asyncio.create_task(coro_trunner(1, 2, 1.2)),
        asyncio.create_task(coro_trunner(2, 3, 2.2)),
        asyncio.create_task(coro_trunner(5, 5, 5.2))
    ]

    try:
        gaf = asyncio.gather(*tasks)
        results = await asyncio.wait_for(gaf,
                                         total_timeout)
    except (asyncio.TimeoutError,
            asyncio.CancelledError,
            Exception):
        # Total timeout reached: get the results that are ready

        # Consume the gather exception
        exc = gaf.exception()  # noqa

        results = []

        for task in tasks:
            if task.done():
                results.append(task.result())
            else:
                # We want to know when a task yields nothing
                results.append(None)
                task.cancel()

    return results


print(f'Total timeout: {total_timeout}')


loop = asyncio.get_event_loop()
start_time = loop.time()

results = asyncio.run(main())
end_time = loop.time()

print(f'{end_time - start_time} --> {results}')

【讨论】:

  • 在 async def work_coro 中,await asyncio.sleep(wait_time) 行实际上应该是 await asyncio.sleep(random_time),不是吗?
  • 您的方法的问题是,如果 task3 超时,您最终会出现“main() 超时”,结果一无所获,即使 task1 和 task2 成功完成。我将 total_timeout 设置为 5 并获得了 (1, 2, 1.2) 的单次生成和 (1, 3, 2.2) 的单次生成,这意味着它们在第一次运行中完成,但是由于 task3 超时,我什么也没得到主要的()。我想得到所有已经完成的东西,而其他所有东西都没有。
  • 确实,我很高,我刚刚按照您的描述进行了工作,稍后我会发布编辑,谢谢
  • 编辑帖子:新版本可以在命令行超时运行(否则为 5 秒)。它适用于 python 3.7、3.8,不知道为什么,但在 3.9 中,主协程不会触发 TimeoutError。
  • 我正在使用 python 3.9.14。在您的代码中,我添加了: import datetime as dt, start_time = dt.datetime.now() (在尝试之前:您设置循环的位置和 print(f"lasted: {round((dt.datetime.now()- start_time).total_seconds(),1)}") 在最后一个 if 之前。对于 5 秒的默认超时,脚本持续了 17.6 秒,这与 5 秒限制相去甚远。在三次尝试中,我总是得到 17.6 秒。在 python 3.10.4我得到了时间:17.6、13.3、17.6。
猜你喜欢
  • 2016-11-07
  • 1970-01-01
  • 2016-10-10
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2014-03-19
  • 1970-01-01
相关资源
最近更新 更多