You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用Python异步并发执行阻塞函数(如requests)

生成持续异步并行的请求流

我有一个会阻塞事件循环的函数(比如发起API请求的函数),需要生成持续的请求流,让这些请求并行执行但非同步启动——也就是每个后续请求要在之前的请求发起之后、完成之前就启动,而非一次性批量创建所有任务再并发执行。

初始尝试方案

我先用loop.run_in_executer()做了尝试,代码如下:

import asyncio
import requests

# blocking_request_func() 已在别处定义

async def main():
    loop = asyncio.get_event_loop()
    future1 = loop.run_in_executor(None, blocking_request_func, 'param')
    future2 = loop.run_in_executor(None, blocking_request_func, 'param')
    response1 = await future1
    response2 = await future2
    print(response1)
    print(response2)

loop = asyncio.get_event_loop()
loop.run_until_complete(main())

这个方案能实现请求并行,但不符合需求:它是先创建一组任务/未来对象,再同步执行这组任务,不是我要的流式启动。

预期执行流程

我需要的执行流程是这样的:

1. 发送request_1,不等待完成。
(在步骤1之后而非同时):
2. 发送request_2,不等待完成。
(在步骤2之后而非同时):
3. 发送request_3,不等待完成。
(request_1或其他请求返回响应)
(在步骤3之后而非同时):
4. 发送request_4,不等待完成。
(request_2或其他请求返回响应)
……以此类推

TaskGroup尝试方案

我又试了asyncio.TaskGroup(),代码如下:

async def request_func():
    global result # 结果列表已在全局区域定义
    loop = asyncio.get_event_loop()
    result.append(await loop.run_in_executor(None, blocking_request_func, 'param'))
    await asyncio.sleep(0) # 加不加这行结果都一样

async def main():
    async with asyncio.TaskGroup() as tg:
       for i in range(0, 10):
           tg.create_task(request_func())

但结果还是先批量创建所有任务,再同步并发执行,没达到流式启动的效果。

接近需求的解决方案(带限制)

后来找到一个最接近需求的方案,但需要手动控制几个参数:

import asyncio
import random
import time
import concurrent.futures

def blockme(n):
    x = random.random() * 2.0
    time.sleep(x)
    return n, x

def cb(fut):
    print("Result", fut.result())
    
async def main():
    # 需要手动控制线程数量
    pool = concurrent.futures.ThreadPoolExecutor(max_workers=4)

    loop = asyncio.get_event_loop()
    futs = []

    # 需要手动控制每秒请求数
    delay = 0.5
    n = 0

    while True:
        fut = loop.run_in_executor(pool, blockme, n)
        fut.add_done_callback(cb)
        futs.append(fut)
        n += 1

        # 需要手动控制同时存在的未完成任务数量,避免堆积过多
        if len(futs) > 40:
            completed, futs = await asyncio.wait(futs, 
                                                 timeout=5, 
                                                 return_when=asyncio.FIRST_COMPLETED)

asyncio.run(main())

这个方案的限制点:

  • 必须手动控制线程池的最大工作线程数
  • 必须手动控制请求的发送间隔(每秒请求数)
  • 必须手动控制同时存在的未完成任务数量,避免任务堆积过多

内容的提问来源于stack exchange,提问作者Alexey Trukhanov

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.03 21:35:18