如何优化基于Python Asyncio的批量HTTP请求程序速度?
嘿,你已经搭好了异步请求的基础框架,还通过实验找到了Semaphore的最优值,这已经很棒了!针对10000次API请求的场景,我整理了几个能进一步拉满速度的优化方向,咱们一个个拆解:
1. 把阻塞式睡眠换成异步睡眠
你的processData里用了time.sleep(0.001),这可是异步编程的大忌——它会直接阻塞整个事件循环,哪怕只有1ms,在10000次请求的累积下,浪费的时间会非常可观。赶紧换成await asyncio.sleep(0.001),这样事件循环在等待的时候还能去处理其他请求,完全不浪费资源。
修改后的代码:
async def processData(data): await asyncio.sleep(0.001) return data
2. 优化ClientSession的连接池配置
aiohttp的ClientSession自带连接池,但默认参数可能没适配你的并发量。你已经把Semaphore设为100,那一定要让连接池的limit和这个值匹配,同时开启长连接,避免重复建立TCP连接的开销:
from aiohttp import TCPConnector async def main(): # 配置连接池,和Semaphore上限一致,开启长连接 connector = TCPConnector(limit=100, keepalive_timeout=30, force_close=False) async with ClientSession(connector=connector) as session: # 你的请求逻辑在这里
3. 批量处理数据,减少上下文切换开销
如果processData的逻辑允许,别拿到一个响应就处理一个,而是先批量获取所有响应数据,再统一处理——这样能减少异步任务的上下文切换次数,提升整体效率。比如:
async def fetch(session, url): async with session.get(url) as response: return await response.json() async def main(): urls = [f"https://api.example.com/{i}" for i in range(10000)] connector = TCPConnector(limit=100) async with ClientSession(connector=connector) as session: # 批量创建请求任务,一次性获取所有数据 tasks = [fetch(session, url) for url in urls] all_data = await asyncio.gather(*tasks) # 批量处理数据(如果processData是同步的,也可以用列表推导式快速处理) processed_data = [processData(data) for data in all_data]
要是processData本身也可以改成异步的,那就用await asyncio.gather(*[processData(d) for d in all_data])来批量处理,效率更高。
4. 用更快的JSON解析库
aiohttp默认的response.json()用的是标准库的json模块,速度不算顶尖。如果API返回的是JSON数据,换成ujson这种更快的解析库,能节省不少解析时间:
import ujson async def fetch(session, url): async with session.get(url) as response: # 先获取文本,再用ujson解析 return ujson.loads(await response.text())
5. 处理CPU密集型任务的特殊方案
如果你的processData涉及大量CPU计算(不只是简单的sleep),那单线程的asyncio会被卡住,因为它没法同时处理CPU任务和IO任务。这时候可以用concurrent.futures.ProcessPoolExecutor把CPU密集任务放到独立进程里处理,让事件循环专注于IO请求:
import concurrent.futures # 假设这是CPU密集的处理函数 def processData(data): # 比如复杂的计算、数据转换 time.sleep(0.001) return data async def main(): # ... 前面的请求逻辑,拿到all_data ... # 用进程池处理CPU任务 with concurrent.futures.ProcessPoolExecutor() as executor: processed_data = await asyncio.get_event_loop().run_in_executor( executor, lambda: [processData(d) for d in all_data] )
6. 优化DNS解析(如果需要)
如果你的API域名需要频繁解析,默认的DNS解析可能成为瓶颈。可以用aiodns这个异步DNS解析库来加速,配置起来也很简单:
import aiodns from aiohttp import TCPConnector async def main(): resolver = aiodns.DNSResolver() connector = TCPConnector(resolver=resolver, limit=100) async with ClientSession(connector=connector) as session: # 你的请求逻辑
7. 用任务队列实现更灵活的调度(可选)
如果需要处理失败重试、分批次执行等复杂逻辑,用asyncio.Queue来做任务调度会更灵活。比如创建固定数量的worker,从队列里取任务执行,这样能更精准地控制并发,还能轻松实现重试:
async def worker(session, queue): while True: url = await queue.get() try: data = await fetch(session, url) await processData(data) except Exception as e: # 简单的重试逻辑,把失败的任务放回队列 print(f"请求失败,重试: {e}") await queue.put(url) finally: queue.task_done() async def main(): queue = asyncio.Queue(maxsize=100) urls = [f"https://api.example.com/{i}" for i in range(10000)] # 把所有URL放入队列 for url in urls: await queue.put(url) connector = TCPConnector(limit=100) async with ClientSession(connector=connector) as session: # 创建100个worker,和Semaphore上限一致 workers = [asyncio.create_task(worker(session, queue)) for _ in range(100)] # 等待所有任务完成 await queue.join() # 取消所有worker任务 for worker_task in workers: worker_task.cancel() await asyncio.gather(*workers, return_exceptions=True)
最后再提个小建议:可以用asyncio.run()来启动你的主函数,比手动管理事件循环更简洁规范。比如:
if __name__ == "__main__": asyncio.run(main())
内容的提问来源于stack exchange,提问作者rajendra

