前言

传统多线程爬虫依托操作系统线程调度实现并发,虽然相比串行爬虫效率显著提升,但线程存在系统资源开销、线程数量上限、线程切换损耗等短板。当面对上千条链接批量采集、海量图片下载场景,多线程依旧存在性能天花板。协程(Coroutine)基于用户态实现并发,由程序自身调度任务,无需操作系统切换线程,内存占用极低、可同时维持数千并发连接,是当前 IO 密集型爬虫最优高性能方案。

Python 生态中 aiohttp 是异步 HTTP 请求标准库,搭配 asyncio 内置协程框架,构建异步爬虫体系。协程非常适合爬虫这类大量网络等待的 IO 任务:当一个协程发起网络请求进入等待状态时,调度器自动切换执行其他就绪协程,充分利用空闲等待时间,在单机实现数百乃至上千并发连接。

协程开发存在学习门槛,异步 / 同步代码不能随意混用、事件循环管理、并发数量限流、信号量控制、异常捕获、会话复用、文件异步写入都是极易踩坑的关键点。无限制创建协程会瞬间建立大量连接,触发目标服务器防火墙、连接拒绝、IP 封禁,因此必须使用信号量限制并发峰值。

本文系统讲解 asyncio 协程基础、aiohttp 使用规范、信号量限流、异步页面采集、异步图片下载、异步文件存储、同步代码兼容方案、高频故障排查,提供工程化可直接运行代码,拆解底层执行原理。

本文核心依赖库官方文档地址:

  1. asyncio Python 内置协程事件循环库:https://docs.python.org/3/library/asyncio.html
  2. aiohttp 异步 HTTP 客户端框架:https://docs.aiohttp.org/
  3. lxml HTML 解析库(同步解析,协程内可直接调用):https://lxml.de/
  4. pathlib 路径管理内置库:https://docs.python.org/3/library/pathlib.html

依赖安装命令:

pip install aiohttp lxml

一、协程爬虫核心基础理论

1.1 协程、线程的核心差异

爬虫属于 IO 密集任务,下表清晰区分方案选型边界:

表格

并发方案 调度主体 资源开销 最大并发规模 适用场景
串行单线程 程序顺序执行 最低 1 少量链接调试、学习测试
threading 多线程 操作系统内核 中等 50~200 中小型常规采集任务
asyncio+aiohttp 协程 程序事件循环调度 极低 500~3000 大规模批量采集、图片批量下载

核心重点:协程不是线程,同一时刻仅有一段代码在 CPU 执行;依靠 IO 等待间隙切换任务,没有操作系统线程切换开销。

1.2 核心关键字语法

协程语法是区分同步代码最直观标识:

  1. async def:定义协程函数,调用不会直接执行,返回协程对象
  2. await:暂停当前协程,等待异步 IO 操作完成,期间释放调度器执行其他任务
  3. async with:异步上下文管理器,用于管理异步会话、异步连接
  4. async for:异步迭代器,爬虫场景使用较少

重要禁忌:await 只能写在 async def 函数内部;同步阻塞函数(requests、time.sleep)不能直接放入协程,会阻塞整个事件循环。

1.3 限流核心:Semaphore 信号量

协程可以轻松创建上千任务,一次性全部发起请求会造成短时间海量访问,直接被服务器封禁。 asyncio.Semaphore 信号量用来控制同时运行的协程数量,等同于多线程的最大并发数,是异步爬虫强制组件。

工作原理:信号量内置计数器,协程执行前获取信号量,计数器减 1;任务结束释放信号量,计数器加 1;达到上限时,后续协程进入等待队列。

1.4 aiohttp 核心组件说明

  1. aiohttp.ClientSession:异步会话对象,类比 requests.Session,复用 TCP 连接,全局建议复用同一个 Session,不要频繁创建销毁。
  2. ClientTimeout:统一设置请求超时时间,防止协程永久阻塞。
  3. 响应对象:支持text()read()(二进制数据,用于图片下载)异步读取。

二、最简异步请求入门案例

先搭建最小可运行示例,理解协程执行流程

import asyncio
import aiohttp

# 定义协程请求函数
async def fetch_page(session: aiohttp.ClientSession, url: str):
    headers = {
        "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36"
    }
    try:
        async with session.get(url, headers=headers, timeout=aiohttp.ClientTimeout(total=10)) as resp:
            # await 异步读取网页文本
            html = await resp.text(encoding="utf-8")
            print(f"请求成功 {url} 页面长度:{len(html)}")
            return html
    except Exception as e:
        print(f"请求失败 {url} 错误:{str(e)}")
        return None

async def main():
    # 创建全局异步会话
    async with aiohttp.ClientSession() as session:
        target_url = "https://www.example.com"
        await fetch_page(session, target_url)

if __name__ == "__main__":
    # 启动事件循环执行主协程
    asyncio.run(main())

代码原理详解

  1. async with aiohttp.ClientSession() 创建会话,程序结束自动关闭连接,释放资源;
  2. session.get() 返回异步请求上下文,async with 自动管理连接释放;
  3. await resp.text() 异步读取响应内容,不会阻塞事件循环;
  4. asyncio.run(main()) Python3.7 + 标准启动方式,自动创建、关闭事件循环。

三、带并发限流的异步列表页爬虫实战

实现分页 URL 批量采集、信号量限流、网页请求、lxml 解析、数据收集完整流程,是资讯爬虫标准模板。

import asyncio
import aiohttp
from lxml import etree

# 全局配置
MAX_CONCURRENT = 15  # 最大并发协程数量(信号量阈值)
PAGE_START = 1
PAGE_END = 30
TIMEOUT = aiohttp.ClientTimeout(total=12)
HEADERS = {
    "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 Chrome/120.0.0.0 Safari/537.36"
}

# 信号量实例
sem = asyncio.Semaphore(MAX_CONCURRENT)

def parse_html(html_text):
    """解析函数为同步函数,协程内可以直接调用"""
    tree = etree.HTML(html_text)
    items = []
    node_list = tree.xpath('//div[@class="news-item"]')
    for node in node_list:
        title = node.xpath('./h3/a/text()')
        link = node.xpath('./h3/a/@href')
        time_info = node.xpath('./span/time/text()')
        data = {
            "title": title[0].strip() if title else "",
            "url": link[0] if link else "",
            "pub_time": time_info[0].strip() if time_info else ""
        }
        items.append(data)
    return items

async def crawl_page(session: aiohttp.ClientSession, page_url: str):
    # 获取信号量,控制并发
    async with sem:
        try:
            async with session.get(page_url, headers=HEADERS, timeout=TIMEOUT) as resp:
                if resp.status != 200:
                    print(f"页面异常状态码 {resp.status}:{page_url}")
                    return []
                html = await resp.text(encoding="utf-8")
                result_data = parse_html(html)
                print(f"成功采集 {page_url},获取{len(result_data)}条资讯")
                return result_data
        except Exception as err:
            print(f"采集失败 {page_url} 错误:{str(err)}")
            return []

async def main():
    # 全局会话
    async with aiohttp.ClientSession() as session:
        task_list = []
        # 批量生成分页任务
        for page in range(PAGE_START, PAGE_END + 1):
            url = f"https://www.example.com/news?page={page}"
            task = asyncio.create_task(crawl_page(session, url))
            task_list.append(task)
        # 等待所有任务完成
        all_results = await asyncio.gather(*task_list)
        # 汇总所有页面数据
        all_news = []
        for page_data in all_results:
            all_news.extend(page_data)
        print(f"全部任务执行完毕,总计采集资讯 {len(all_news)} 条")
        # 数据落地写入文件
        with open("./async_crawl_result.txt", "w", encoding="utf-8") as f:
            for item in all_news:
                line = f"{item['title']} | {item['url']} | {item['pub_time']}\n"
                f.write(line)

if __name__ == "__main__":
    asyncio.run(main())

核心原理拆解

  1. asyncio.Semaphore 放置在请求函数内部,每个任务执行前抢占信号量,严格限制同时活跃请求数量;
  2. asyncio.create_task() 创建协程任务,提交给事件循环调度;
  3. asyncio.gather(*task_list) 批量等待所有任务完成,收集所有返回结果;
  4. 数据解析使用同步 lxml,少量 CPU 计算不会阻塞事件循环;如果存在重度计算,建议使用线程池执行解析;
  5. 文件写入统一在所有任务结束后同步执行,避免异步文件 IO 锁竞争问题。

四、异步图片批量下载实战(二进制流读取)

图片下载需要读取二进制响应内容,使用await resp.read(),结合信号量限流,实现高性能图片批量采集。

import asyncio
import aiohttp
from pathlib import Path

MAX_DOWNLOAD_NUM = 12
sem = asyncio.Semaphore(MAX_DOWNLOAD_NUM)
SAVE_DIR = Path("./async_images")
SAVE_DIR.mkdir(exist_ok=True, parents=True)
TIMEOUT = aiohttp.ClientTimeout(total=15)
HEADERS = {
    "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36",
    "Referer": "https://www.example.com/"
}

async def download_img(session: aiohttp.ClientSession, img_url: str, save_name: str):
    async with sem:
        try:
            async with session.get(img_url, headers=HEADERS, timeout=TIMEOUT) as resp:
                if resp.status != 200:
                    print(f"图片访问失败 {img_url} status:{resp.status}")
                    return False
                # 异步读取二进制数据流
                img_bytes = await resp.read()
                save_path = SAVE_DIR / save_name
                # 文件写入为同步操作,短时IO无压力
                with open(save_path, "wb") as f:
                    f.write(img_bytes)
                print(f"下载完成:{save_name}")
                return True
        except Exception as err:
            print(f"下载异常 {img_url} 错误:{str(err)}")
            return False

async def main():
    img_url_list = [
        "https://example.com/img/1.jpg",
        "https://example.com/img/2.png",
        "https://example.com/img/3.webp"
    ]
    async with aiohttp.ClientSession() as session:
        task_list = []
        for idx, url in enumerate(img_url_list):
            suffix = url.split(".")[-1].split("?")[0]
            filename = f"async_img_{idx+1:03d}.{suffix}"
            task = asyncio.create_task(download_img(session, url, filename))
            task_list.append(task)
        await asyncio.gather(*task_list)
    print("所有图片下载任务结束")

if __name__ == "__main__":
    asyncio.run(main())

关键要点

  1. 图片二进制数据采用await resp.read()获取,禁止使用 text ();
  2. 大批量高清图片场景,可以分块异步读取响应;
  3. 文件名生成逻辑提前处理,避免协程内部并发统计文件数量产生竞争;
  4. 添加 Referer 请求头绕过基础图片防盗链。

五、协程爬虫高频致命误区

5.1 误区 1:协程内部使用 requests、time.sleep

requests 是同步阻塞库,一旦调用,整个事件循环全部冻结,所有协程暂停调度。

  • 替换方案:网络请求统一使用 aiohttp;延时使用await asyncio.sleep(0.5)

5.2 误区 2:大量重复创建 ClientSession

ClientSession 内部维护连接池,频繁创建销毁会极大损耗性能,工程规范:全局仅创建一个 Session。

5.3 误区 3:不使用 Semaphore,无限制并发

直接一次性创建上千任务,瞬间发起海量连接,目标服务器直接封禁 IP,本机也会出现连接创建失败。所有生产环境异步爬虫必须配置信号量限流。

5.4 误区 4:忽略异常捕获,单个任务崩溃导致全部任务终止

协程内部未捕获异常时,不会直接崩溃程序,但容易出现任务静默失败,务必使用 try-except 包裹请求逻辑。

5.5 误区 5:混用同步异步文件读写

大量协程同时打开写入同一个文件,极易出现文件内容错乱。最优方案:协程只采集数据,全部任务结束后统一写入文件。

六、同步代码调用协程、协程调用同步阻塞函数方案

6.1 在普通同步函数运行协程

# 同步函数内部启动协程
def run_async_task():
    async def demo():
        await asyncio.sleep(0.1)
        print("异步任务执行")
    asyncio.run(demo())

6.2 协程中运行阻塞同步函数(防止阻塞事件循环)

使用asyncio.to_thread()把同步函数放入独立线程运行,不阻塞事件循环:

def heavy_sync_parse(html):
    # 重度同步解析代码
    pass

async def work():
    html_data = await fetch()
    # 同步函数放入线程执行
    result = await asyncio.to_thread(heavy_sync_parse, html_data)

七、多线程爬虫 VS 协程爬虫选型标准

  1. 小规模采集(并发 < 30):多线程 threading 开发简单,足够使用;
  2. 大规模批量采集(并发 50~2000):优先 aiohttp 协程,内存占用更低、吞吐更高;
  3. 存在大量 CPU 计算任务:协程优势消失,改用多进程;
  4. Windows 旧 Python 版本、老旧服务器环境:优先多线程,部分环境协程事件循环存在兼容性问题。

八、工程化拓展优化方向

  1. 异步请求代理 IP 支持:aiohttp 支持配置代理,实现异步分布式 IP 轮换;
  2. 异步队列解耦任务:搭配asyncio.Queue实现异步版本生产者消费者模型;
  3. 请求频率控制:在协程内部增加await asyncio.sleep()随机延时,平滑访问速率;
  4. 会话 Cookie 持久化:封装异步会话,实现登录态保持,模拟登录采集;
  5. 限制并发之外的熔断机制:连续多次失败链接自动加入黑名单,避免无效请求。

九、总结

协程异步爬虫是单机高性能采集的最优方案,核心要点汇总:

  1. 使用async def定义协程函数,IO 操作使用await切换任务;
  2. 全局复用aiohttp.ClientSession,提升连接复用率;
  3. 必须使用 Semaphore 信号量限制最大并发,杜绝无限并发;
  4. 禁止在协程内部调用同步阻塞库 requests、time.sleep;
  5. asyncio.gather批量调度任务,统一收集结果;
  6. 数据落地尽量在所有异步任务完成后集中写入,规避文件竞争问题。
Logo

openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构

更多推荐