【问题标题】:Python Aiohttp Asyncio: how to create delays between each taskPython Aiohttp Asyncio:如何在每个任务之间创建延迟
【发布时间】:2022-11-29 05:48:58
【问题描述】:

我试图解决的问题:我正在向服务器发出许多 api 请求。我试图在异步 api 调用之间创建延迟以符合服务器的速率限制策略。

我想让它做什么我希望它表现得像这样:

  1. 发出 api 请求#1
  2. 等待0.1秒
  3. 发出 api 请求 #2
  4. 等待0.1秒 ... 等等 ...
  5. 重复直到发出所有请求
  6. 收集响应并将结果返回一个对象(结果)

    问题:我什么时候介绍的asyncio.sleep()或者时间.睡眠()在代码中,它仍然几乎是瞬间发出 api 请求。它似乎延迟了执行打印(),但不是 api 请求。我怀疑我必须在环形,而不是在 fetch_one() 或 fetch_all() 处,但不知道该怎么做。

    代码块:

    async def fetch_all(loop, urls, delay): 
        results = await asyncio.gather(*[fetch_one(loop, url, delay) for url in urls], return_exceptions=True)
        return results
    
    async def fetch_one(loop, url, delay):
    
        #time.sleep(delay)
        #asyncio.sleep(delay)
    
        async with aiohttp.ClientSession(loop=loop) as session:
            async with session.get(url, ssl=SSLContext()) as resp:
                # print("An api call to ", url, " is made at ", time.time())
                # print(resp)
                return await resp
    
    delay = 0.1
    urls = ['some string list of urls']
    loop = asyncio.get_event_loop()
    loop.run_until_complete(fetch_all(loop, urls, delay))
    
    Versions I'm using: 
    python                    3.8.5
    aiohttp                   3.7.4
    asyncio                   3.4.3
    

    我将不胜感激指导我走向正确方向的任何提示!

【问题讨论】:

    标签: python api python-asyncio aiohttp delayed-execution


    【解决方案1】:

    asyncio.gather 的调用将“同时”启动所有请求——另一方面,如果您只是对每个任务使用锁或等待,那么您根本不会从使用并行性中获得任何好处。

    最简单的做法是,如果您知道可以发出请求的速率,只需在每个连续请求之前增加异步暂停时间——一个简单的全局变量可以做到这一点:

    
    next_delay = 0.1
    
    async def fetch_all(loop, urls, delay): 
        results = await asyncio.gather(*[fetch_one(loop, url, delay) for url in urls], return_exceptions=True)
        return results
    
    async def fetch_one(loop, url, delay):
        global next_delay
        
        next_delay += delay
        await asyncio.sleep(next_delay)
    
        async with aiohttp.ClientSession(loop=loop) as session:
            async with session.get(url, ssl=SSLContext()) as resp:
                # print("An api call to ", url, " is made at ", time.time())
                # print(resp)
                return await resp
    
    delay = 0.1
    urls = ['some string list of urls']
    loop = asyncio.get_event_loop()
    loop.run_until_complete(fetch_all(loop, urls, delay))
    

    现在,如果你想发出 5 个请求,然后发出下 5 个请求,你可以使用像 asyncio.Condition 这样的同步原语,在检查有多少 api 调用处于活动状态的表达式上使用它的 wait_for

    active_calls = 0
    
    MAX_CALLS = 5
    
    async def fetch_all(loop, urls, delay): 
        event = asyncio.Event()
        event.set()
        results = await asyncio.gather(*[fetch_one(loop, url, delay, event) for url in urls], return_exceptions=True)
        return results
    
    async def fetch_one(loop, url, delay, cond):
        global active_calls
        
        active_calls += 1
        if active_calls > MAX_CALLS:
            event.clear()
            
        await event.wait()
        
        try:
            async with aiohttp.ClientSession(loop=loop) as session:
                async with session.get(url, ssl=SSLContext()) as resp:
                    # print("An api call to ", url, " is made at ", time.time())
                    # print(resp)
                    return await resp
        finally:
            active_calls -= 1
        if active_calls == 0:
            event.set()
            
    
    urls = ['some string list of urls']
    loop = asyncio.get_event_loop()
    loop.run_until_complete(fetch_all(loop, urls, delay))
    

    对于这两个示例,如果您的任务在设计中避免使用全局变量(实际上,这些是“模块”变量) - 您可以将所有功能移动到一个类中,并在一个实例上工作,并将全局变量提升为实例属性,或者使用可变容器,例如在其第一项中保存 active_calls 值的列表,并将其作为参数传递。

    【讨论】:

    • 谢谢你给我这么好的提示!在 Python 中做事的工具和方法太多了。
    • 我到底在找什么,谢谢
    • 很好的答案,谢谢!但是为什么我们(在你的第一个例子中)递增next_delay?直觉上,我希望这意味着最后一个,比如说,1000 个请求休眠 100 秒而不是 0.1 秒,但不知何故,情况似乎并非如此。
    • next_delay 增加是为了让每个任务将其“核心”部分的执行间隔开来,因此将实际网络请求间隔为 delay 的值。它更像是一个示例,因为这会阻止任何实际的并行性并一次发出一个请求。稍微调整一下,就可以更改为一次并行发出 10 个请求,从而保持 api 使用吞吐量(实际上,只需微调此处的“延迟”量就可以做到这一点。
    【解决方案2】:

    当您使用asyncio.gather 时,您同时运行所有fetch_one 协程。他们一起等待delay,而不是一起即时调用 API。

    要解决这个问题,您应该在fetch_all 中一个一个地等待fetch_one,或者使用Semaphore 来表示下一个不应该在上一个完成之前开始。

    这是想法:

    import asyncio
    
    _sem = asyncio.Semaphore(1)
    
    
    async def fetch_all(loop, urls, delay): 
        results = await asyncio.gather(*[fetch_one(loop, url, delay) for url in urls], return_exceptions=True)
        return results
    
    async def fetch_one(loop, url, delay):
    
        async with _sem:  # next coroutine(s) will stuck here until the previous is done
            await asyncio.sleep(delay)
    
            async with aiohttp.ClientSession(loop=loop) as session:
                async with session.get(url, ssl=SSLContext()) as resp:
                    # print("An api call to ", url, " is made at ", time.time())
                    # print(resp)
                    return await resp
    
    delay = 0.1
    urls = ['some string list of urls']
    loop = asyncio.get_event_loop()
    loop.run_until_complete(fetch_all(loop, urls, delay))
    

    【讨论】:

    • 您使用asyncio.Semaphore 而不是asyncio.Lock 是否有特殊原因?
    • @ŁukaszKwieciński 没有特别的理由,Lock 也可以在这里使用。我会说信号量可能是轻微地如果将来 OP 希望一次允许多个请求,那就更好了。
    猜你喜欢
    • 1970-01-01
    • 2020-07-15
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-09-28
    • 1970-01-01
    • 1970-01-01
    • 2020-12-12
    相关资源
    最近更新 更多