如何用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
相关产品推荐
相关产品推荐

