Python asyncio任务无法运行及并发控制方案正确性咨询
Asyncio与aiohttp并行抓取API数据的问题排查与实现验证
问题背景
测试Python的asyncio和aiohttp库,需求是并行从API抓取数据并将.html文件保存到本地磁盘。初始代码抛出两类错误:
- TypeError: Passing coroutines is forbidden, use tasks explicitly.
- RuntimeWarning: coroutine 'web_scrape_task' was never awaited
之后尝试用asyncio.Semaphore将并发调用限制为2并修改了代码,需验证该实现是否正确。
初始代码
import asyncio import aiohttp import time import os url_i = "<some_urls>-" file_path = "<local_drive>\\asynciotest" async def download_pep(pep_number: int) -> bytes: # 原代码此处误用未定义变量url,应改为url_i url = url_i + f"{pep_number}/" print(f"Begin downloading {url}") async with aiohttp.ClientSession() as session: async with session.get(url) as resp: content = await resp.read() print(f"Finished downloading {url}") return content async def write_to_file(pep_number: int, content: bytes) -> None: with open(os.path.join(file_path,f"{pep_number}"+'-async.html'), "wb") as pep_file: print(f"{pep_number}_Begin writing ") pep_file.write(content) print(f"Finished writing") async def web_scrape_task(pep_number: int) -> None: content = await download_pep(pep_number) await write_to_file(pep_number, content) async def main() -> None: tasks = [] for i in range(8010, 8016): # 错误:直接添加协程对象到列表,未显式创建任务 tasks.append(web_scrape_task(i)) await asyncio.wait(tasks) if __name__ == "__main__": s = time.perf_counter() asyncio.run(main()) elapsed = time.perf_counter() - s print(f"Execution time: {elapsed:0.2f} seconds.")
初始代码报错原因
- TypeError:
asyncio.wait()要求传入任务(Task)对象而非协程对象,原代码直接将web_scrape_task(i)(协程)加入列表,未用asyncio.create_task()包装。 - RuntimeWarning:属于TypeError引发的连锁问题,协程未被正确调度执行。
修改后代码的问题
修改后的代码存在两处关键错误,完全无法实现预期的并行+并发限制效果:
async def web_scrape_task(pep_number: int) -> None: # 错误:每个任务内部循环处理所有编号,导致重复抓取 for i in range(8010, 8016): content = await download_pep(i) # 错误:写入时用外层的pep_number而非当前循环的i,所有内容写入同一文件 await write_to_file(pep_number, content) ##To limit concurrent call 2## sem = asyncio.Semaphore(2) async def main() -> None: tasks = [] for i in range(8010, 8016): # 错误:循环内直接用async with sem,导致串行创建任务,完全失去并行性 async with sem: tasks.append(asyncio.create_task(web_scrape_task(i))) await asyncio.gather(*tasks) if __name__ == "__main__": s = time.perf_counter() asyncio.run(main()) elapsed = time.perf_counter() - s print(f"Execution time: {elapsed:0.2f} seconds.")
正确实现方案
核心修正点
- 恢复
web_scrape_task的单一职责:每个任务只处理一个编号,避免内部循环。 - 将
Semaphore的获取放在网络请求的协程内部,确保并发限制作用于实际的API调用。 - 正确创建任务并使用
asyncio.gather()调度。
正确代码示例
import asyncio import aiohttp import time import os url_i = "<some_urls>-" file_path = "<local_drive>\\asynciotest" # 定义全局信号量,限制并发数为2 sem = asyncio.Semaphore(2) async def download_pep(pep_number: int) -> bytes: url = url_i + f"{pep_number}/" print(f"Begin downloading {url}") # 在网络请求前获取信号量,限制并发 async with sem: async with aiohttp.ClientSession() as session: async with session.get(url) as resp: content = await resp.read() print(f"Finished downloading {url}") return content async def write_to_file(pep_number: int, content: bytes) -> None: # 确保目标目录存在,避免写入失败 os.makedirs(file_path, exist_ok=True) file_name = os.path.join(file_path, f"{pep_number}-async.html") with open(file_name, "wb") as pep_file: print(f"{pep_number}_Begin writing ") pep_file.write(content) print(f"{pep_number}_Finished writing") async def web_scrape_task(pep_number: int) -> None: content = await download_pep(pep_number) await write_to_file(pep_number, content) async def main() -> None: tasks = [] for i in range(8010, 8016): # 显式创建任务对象 task = asyncio.create_task(web_scrape_task(i)) tasks.append(task) # 等待所有任务完成 await asyncio.gather(*tasks) if __name__ == "__main__": s = time.perf_counter() asyncio.run(main()) elapsed = time.perf_counter() - s print(f"Execution time: {elapsed:0.2f} seconds.")
代码说明
Semaphore放在download_pep内部的网络请求块前,确保同时最多只有2个API请求在执行。- 每个
web_scrape_task只处理一个编号,避免重复操作。 - 用
asyncio.create_task()包装协程,符合asyncio的调度要求。 - 添加
os.makedirs(file_path, exist_ok=True)确保目标目录存在,避免写入报错。
内容的提问来源于stack exchange,提问作者pythondumb
相关产品推荐
相关产品推荐

