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

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.")

初始代码报错原因

  1. TypeError:asyncio.wait()要求传入任务(Task)对象而非协程对象,原代码直接将web_scrape_task(i)(协程)加入列表,未用asyncio.create_task()包装。
  2. 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.")

正确实现方案

核心修正点

  1. 恢复web_scrape_task的单一职责:每个任务只处理一个编号,避免内部循环。
  2. 将Semaphore的获取放在网络请求的协程内部,确保并发限制作用于实际的API调用。
  3. 正确创建任务并使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 10:09:59