为何使用asyncio.gather处理pandas DataFrame行时无法实现并发执行?
为什么我的aiohttp异步请求和同步速度一致?
你的代码本身的异步逻辑写法是正确的——任务创建与并发等待的流程没有问题,导致请求无法并发执行的核心原因大概率是目标服务器的并发限制,而非代码本身的问题。
具体分析
你使用的http://httpbin.org/delay/0.1端点,虽然是用来模拟延迟请求的测试接口,但httpbin服务器对单个IP的并发请求数量有严格限制(通常是限制为1)。这意味着即使你发起100个异步请求,服务器也会串行处理每一个,每个请求耗时0.1秒,总耗时自然就是100*0.1=10秒,和同步循环的结果完全一致。
另外你怀疑的几个点都不是问题:
df.iterrows()只是同步遍历行数据来创建异步任务,这个过程不会阻塞事件循环;asyncio.create_task会立即把任务加入事件循环队列,不会等待任务执行,循环创建100个任务的开销几乎可以忽略。
验证方法
你可以通过以下方式验证这个结论:
- 减少请求数量:把DataFrame的行数改成10,同步耗时应该是1秒;如果异步能在0.1-0.2秒完成,说明代码的并发逻辑是正常的;
- 搭建本地测试服务器:自己写一个简单的异步延迟服务,就能看到真实的并发效果。示例代码:
from fastapi import FastAPI import asyncio app = FastAPI() @app.get("/delay/{seconds}") async def delay(seconds: float): await asyncio.sleep(seconds) return {"message": "Delayed response"}
启动服务后,用http://localhost:8000/delay/0.1作为请求URL,测试100个请求的异步耗时,就能达到你预期的0.1-0.2秒。
针对50万行数据的优化建议
如果你的实际业务场景不是用httpbin,而是自有服务,建议做以下优化:
- 显式配置连接池:虽然aiohttp默认并发限制是100,但显式设置可以避免环境差异影响:
async with aiohttp.ClientSession(connector=aiohttp.TCPConnector(limit=200)) as session: # 后续逻辑不变
- 添加请求超时:避免个别慢请求拖慢整体处理速度:
async with session.get(url, timeout=aiohttp.ClientTimeout(total=5)) as response: await response.text() return row_id
- 分批次处理任务:50万行一次性创建任务会占用大量内存,建议分批次发起请求:
async def process_batch(session, batch): tasks = [asyncio.create_task(fetch(session, row['url'], row['id'])) for _, row in batch.iterrows()] return await asyncio.gather(*tasks) async def process_dataframe(df, batch_size=200): async with aiohttp.ClientSession(connector=aiohttp.TCPConnector(limit=batch_size)) as session: results = [] for i in range(0, len(df), batch_size): batch = df.iloc[i:i+batch_size] batch_results = await process_batch(session, batch) results.extend(batch_results) return results
内容的提问来源于stack exchange,提问作者Джон Сноу
相关产品推荐
相关产品推荐

