为什么在使用asyncio.gather处理Pandas DataFrame的行时,不能实现并发?
我在处理一个包含50万行的pandas DataFrame,每行都需要一个HTTP请求。我把 requests 换成 aiohttp + asyncio,本以为能显著提速,结果异步版本的速度和同步循环一样慢。我搞不清到底并发被卡在了哪里。
下面是一个最小可复现的示例:
import asyncio
import aiohttp
import pandas as pd
import time
# Sample DataFrame with 100 rows, each hitting a delayed endpoint
df = pd.DataFrame({
'id': range(100),
'url': ['http://httpbin.org/delay/0.1'] * 100
})
async def fetch(session, url, row_id):
async with session.get(url) as response:
await response.text()
return row_id
async def process_dataframe(df):
async with aiohttp.ClientSession() as session:
tasks = []
for _, row in df.iterrows():
task = asyncio.create_task(fetch(session, row['url'], row['id']))
tasks.append(task)
results = await asyncio.gather(*tasks)
return results
start = time.time()
results = asyncio.run(process_dataframe(df))
print(f"Async time: {time.time() - start:.2f} seconds")
预期: 对于一个延迟0.1秒的端点,100次请求时并发应该在0.1-0.2秒内完成。
实际: 脚本耗时10秒,与同步循环完全相同。
我已经排查了以下几点:
- 事件循环正在运行(没有关于缺少循环的
RuntimeError提示)。 aiohttp已正确安装,且没有关于DNS/解析器阻塞的警告。- 在macOS与 Linux上都同样如此,Python 3.11与 3.12均如此。
问题: 为什么请求没有被并发发送?在异步函数中遍历DataFrame的每一行会不会以某种方式把任务串行化?还是在紧密循环中使用 asyncio.create_task 有已知的限制?
环境:
- Python 3.11.7
- pandas 2.1.4
- aiohttp 3.9.1
解决方案
我修改了代码,将网络调用替换为sleep 0.1秒的调用:
async def fetch(session, url, row_id):
await asyncio.sleep(.1)
return row_id
"""
async with session.get(url) as response:
await response.text()
return row_id
"""
程序输出:
Async time: 0.10 seconds
这正是人们所期望的。当我运行你原始的代码时,得到的结果从0.71秒到5.54秒不等。这并不奇怪。也许你没有考虑到向服务器传输数据来回的额外时间,一旦发出请求,服务器端可能会因为资源而产生竞争。我们真的知道服务器是否能高效地处理100个并发请求吗?因此,在理论上总耗时的下限应该是0.1秒,前提是数据可以瞬时传输并即时处理,除了所请求的这0.1秒延迟之外。
总之,asyncio 的表现符合预期,所谓的问题其实出在网络层。
站内所有文章版权归属LeftHeroAI导航站,无授权禁止任何主体转载、抄袭、复制内容,亦不得私自架设镜像站点。一经侵权,本站将通过法律途径追责。