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

如何用asyncio实现ThreadPoolExecutor的有限并行批量请求功能?

用asyncio实现带并行限制的批量请求

问题背景

我使用ThreadPoolExecutor时,可以通过以下方式发送带并行请求限制的批量请求:

from concurrent.futures import ThreadPoolExecutor

MAX_PARALLEL_REQUESTS = 5

def request_func(input_data):
    # 模拟同步请求逻辑
    return f"Result for {input_data}"

requests_input_data = [1, 2, 3, 4, 5, 6, 7]

with ThreadPoolExecutor(max_workers=MAX_PARALLEL_REQUESTS) as pool:
    results = list(pool.map(request_func, requests_input_data))

请问如何用asyncio实现相同行为?是否有相关库可用,还是需自行编写「等待首个Future完成后添加新请求」的逻辑?


实现方案

1. 原生asyncio结合Semaphore控制并发

完全不需要自己写复杂的「等待首个Future完成再添加新请求」逻辑,用asyncio.Semaphore就能轻松限制并发数,配合asyncio.gather即可实现和ThreadPoolExecutor一致的效果:

方式一:信号量包装异步请求(推荐)

import asyncio

MAX_PARALLEL_REQUESTS = 5

async def async_request_func(input_data, semaphore):
    async with semaphore:
        # 替换为实际异步IO逻辑(比如异步HTTP请求、数据库操作)
        await asyncio.sleep(1)
        return f"Result for {input_data}"

async def main():
    semaphore = asyncio.Semaphore(MAX_PARALLEL_REQUESTS)
    requests_input_data = [1, 2, 3, 4, 5, 6, 7]
    
    # 创建所有异步任务,每个任务通过信号量限制并发
    tasks = [async_request_func(data, semaphore) for data in requests_input_data]
    # 等待所有任务完成并收集结果
    results = await asyncio.gather(*tasks)
    print(results)

if __name__ == "__main__":
    asyncio.run(main())

信号量会自动控制同时执行的任务数,超出限制的任务会进入等待队列,直到有任务释放信号量,底层由asyncio调度,无需手动管理任务生命周期。

方式二:分批提交任务(适合超大任务量)

如果任务量极大,一次性创建所有任务可能占用过多内存,可以按并发数拆分批次,逐批执行:

import asyncio

MAX_PARALLEL_REQUESTS = 5

async def async_request_func(input_data):
    await asyncio.sleep(1)
    return f"Result for {input_data}"

async def process_batch(batch):
    tasks = [async_request_func(data) for data in batch]
    return await asyncio.gather(*tasks)

async def main():
    requests_input_data = [1, 2, 3, 4, 5, 6, 7]
    # 按并发数拆分任务批次
    batches = [requests_input_data[i:i+MAX_PARALLEL_REQUESTS] for i in range(0, len(requests_input_data), MAX_PARALLEL_REQUESTS)]
    
    results = []
    for batch in batches:
        batch_results = await process_batch(batch)
        results.extend(batch_results)
    print(results)

if __name__ == "__main__":
    asyncio.run(main())

2. 第三方库简化实现

如果是HTTP请求场景,aiohttp这类成熟库自带连接池控制,无需手动加信号量:

import aiohttp
import asyncio

MAX_PARALLEL_REQUESTS = 5

async def async_request_func(session, url):
    async with session.get(url) as response:
        return await response.text()

async def main():
    urls = ["https://example.com"] * 7  # 模拟请求列表
    # 通过TCPConnector限制并发连接数
    connector = aiohttp.TCPConnector(limit=MAX_PARALLEL_REQUESTS)
    async with aiohttp.ClientSession(connector=connector) as session:
        tasks = [async_request_func(session, url) for url in urls]
        results = await asyncio.gather(*tasks)
        print(results)

if __name__ == "__main__":
    asyncio.run(main())

另外还有asyncio-throttle这类专门做并发/速率限制的库,适合更复杂的调度需求,但原生Semaphore已经能覆盖大多数基础场景。

总结

  • 无需自行编写任务等待逻辑,原生asyncio.Semaphore足以实现并发数限制,用法简单高效;
  • HTTP请求优先选择aiohttp这类带连接池的库,无需额外处理并发控制;
  • 超大任务量可采用分批提交的方式,避免内存占用过高。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 03:22:39