Python异步爬虫进阶:面向对象封装实战(以电影网站为例)

一、前言

在上一篇文章中,我们实现了基础的异步爬虫。本文将代码进行面向对象重构,封装成一个可复用的AsyncScraper类,让代码更模块化、更易维护和扩展。

项目地址:(一个静态电影网站,适合爬虫练习)地址不公开

二、版本对比:过程式 vs 面向对象

2.1 过程式版本的痛点

  • 函数之间耦合度高
  • 配置参数分散
  • 难以复用和扩展
  • 连接池、SSL等管理混乱

2.2 面向对象版本的优势

  • 封装性:请求逻辑、配置、生命周期都封装在类中
  • 可配置性:初始化时可设置各种参数
  • 可复用性:可以创建多个爬虫实例
  • 易于维护:每个方法职责单一

三、完整代码实现

3.1 导入所需模块

import json
import time
import aiohttp
import asyncio
import random
from fake_useragent import UserAgent
from lxml import etree
import logging
import ssl
from typing import Optional, Dict, Tuple, List

3.2 日志和SSL配置函数

def logger_config():
    """配置日志"""
    logging.basicConfig(level=logging.INFO,
                        format='%(asctime)s - %(filename)s - %(levelname)s - %(message)s')
    logger = logging.getLogger()
    return logger

def ssl_config():
    """配置SSL(解决证书验证问题)"""
    ssl_context = ssl.create_default_context()
    ssl_context.check_hostname = False
    ssl_context.verify_mode = ssl.CERT_NONE
    return ssl_context

3.3 核心爬虫类 AsyncScraper

class AsyncScraper:
    def __init__(
        self,
        retry: int = 3,
        time_sleep: Tuple[float, float] = (0.5, 1.5),
        timeout: int = 5,
        headers: Optional[Dict[str, str]] = None,
        semaphore: Optional[int] = 3,
        logger: Optional[logging.Logger] = None,
        ssl_context: Optional[ssl.SSLContext] = None,
        ua: Optional[UserAgent] = None,
        limit_per_host: int = 5,
        limit: int = 10,
    ):
        """
        异步爬虫类的初始化方法
        
        :param retry: 重试次数
        :param time_sleep: 请求间隔范围 (min, max)
        :param timeout: 超时时间
        :param headers: 请求头
        :param semaphore: 协程并发数
        :param logger: 日志实例
        :param ssl_context: SSL配置
        :param ua: UserAgent实例
        :param limit_per_host: 单域名最大连接数
        :param limit: 全局最大连接数
        """
        self.retry = retry
        self.time_sleep = time_sleep
        self.timeout = aiohttp.ClientTimeout(total=timeout)
        self.headers = headers or {}
        self.ua = ua or UserAgent()
        self.logger = logger or logging.getLogger()
        self.session = None

        # 并发控制:信号量
        self.semaphore = asyncio.Semaphore(semaphore)

        # SSL配置
        self.ssl_context = ssl_context or ssl.create_default_context()

        # 连接池配置
        self.connector = aiohttp.TCPConnector(
            limit=limit,
            limit_per_host=limit_per_host,
            ssl=self.ssl_context
        )

3.4 核心请求方法

    # ================= 单请求方法(内部使用) =================
    async def _fetch(self, url: str):
        """
        发送单个HTTP请求,支持重试机制
        
        :param url: 请求URL
        :return: 响应文本或None
        """
        for attempt in range(self.retry):
            try:
                async with self.semaphore:
                    # 动态生成User-Agent
                    current_headers = self.headers.copy()
                    current_headers["User-Agent"] = self.ua.random

                    self.logger.info(f"正在请求: {url}")

                    async with self.session.get(
                        url,
                        headers=current_headers,
                        timeout=self.timeout
                    ) as response:

                        status = response.status

                        # 2xx 成功
                        if 200 <= status < 300:
                            await asyncio.sleep(random.uniform(*self.time_sleep))
                            self.logger.info(f"请求成功: {url}")
                            return await response.text()

                        # 4xx 客户端错误(不重试)
                        elif 400 <= status < 500:
                            self.logger.warning(f"客户端错误 {status}: {url}")
                            return None

                        # 5xx 服务器错误(重试)
                        elif 500 <= status < 600:
                            wait = random.uniform(*self.time_sleep) + attempt * 0.2
                            self.logger.warning(f"服务器错误 {status} 重试 {attempt + 1}")
                            await asyncio.sleep(wait)

            except (aiohttp.ClientError, asyncio.TimeoutError) as e:
                self.logger.warning(f"网络异常 {e} -> 重试 {attempt + 1}")
                if attempt < self.retry - 1:
                    await asyncio.sleep(random.uniform(*self.time_sleep))

            except Exception as e:
                self.logger.error(f"未知错误: {e}")
                return None

        return None

    # ================= 批量请求方法(对外使用) =================
    async def fetch_many(self, urls: List[str]):
        """
        并发批量请求多个URL
        
        :param urls: URL列表
        :return: 响应列表
        """
        tasks = [asyncio.create_task(self._fetch(url)) for url in urls]
        return await asyncio.gather(*tasks)

3.5 声明周期管理

    ####################### 生命周期 ###########################
    async def start(self):
        """启动爬虫,创建会话"""
        self.session = aiohttp.ClientSession(
            connector=self.connector,
            timeout=self.timeout
        )
        self.logger.info('爬虫引擎启动')

    async def close(self):
        """关闭爬虫,释放资源"""
        await self.session.close()
        self.logger.info('爬虫引擎已关闭')

3.6 页面解析函数

def def_next_urls(response_text, logger):
    """
    从列表页提取详情页URL
    
    :param response_text: HTML文本
    :param logger: 日志实例
    :return: URL列表
    """
    if not response_text:
        logger.error("def_next_urls 获得的 response 为空")
        return []

    html = etree.HTML(response_text)
    if html is None:
        logger.error("HTML解析失败")
        return []

    next_xpath_url = html.xpath('//a[@class="name"]/@href')
    return next_xpath_url


def dict_data(response_text, logger) -> dict:
    """
    从详情页提取电影信息
    
    :param response_text: HTML文本
    :param logger: 日志实例
    :return: 电影信息字典
    """
    if not response_text:
        return {'返回结果': None, 'RESPONST': 'data函数中传入的response为空'}
    
    html = etree.HTML(response_text)
    if html is None:
        logger.error("HTML解析失败")
        return {}

    try:
        # 提取电影名称
        name = html.xpath('//h2[@class="m-b-sm"]/text()')
        
        # 提取电影分类
        categories = html.xpath('//*[contains(@class, "categories")]'
                                '//button[contains(@class, "el-button")]//span/text()')

        # 提取基本信息(地区、时长)
        info_spans = html.xpath("//div[@class='m-v-sm info']/span/text()")

        # 提取评分
        score = html.xpath('//p[@class="score m-t-md m-b-n-sm"]/text()')

        # 提取剧情简介
        drama = html.xpath('//*[contains(@class, "drama")]/p[1]/text()')

        detail_dict = {
            'name': name[0] if name else None,
            'categories': categories if categories else None,
            'categories_str': ','.join([cat.strip() for cat in categories] if categories else []),
            'region': info_spans[0].strip() if len(info_spans) > 0 else '未知',
            'duration': info_spans[2].strip() if len(info_spans) > 2 else '未知',
            'score': score[0].strip() if score else '0.0',
            'drama': drama[0].strip() if drama else '',
        }
        logger.info(f'成功解析电影: {detail_dict["name"]}')

        return detail_dict

    except Exception as e:
        logger.warning(f'解析电影信息时出错: {e}')
        return {'解析电影时出错': e}

3.7 主函数

async def main():
    """主函数:协调整个爬虫流程"""
    
    # 1. 配置初始化
    logger = logger_config()
    ssl_context = ssl_config()
    ua = UserAgent()
    
    # 2. 请求配置
    url = '不公开,怕被揍'
    headers = {
        'User-Agent': None,  # 动态生成
        'Accept': 'text/html,application/xhtml+xml,application/xml;q=0.9,*/*;q=0.8',
        'Accept-Language': 'zh-CN,zh;q=0.8,en-US;q=0.5,en;q=0.3',
        'Accept-Encoding': 'gzip, deflate',
        'Connection': 'keep-alive',
    }
    
    # 3. 生成列表页URL(共10页)
    index_urls = [url + '/page/' + str(i) for i in range(1, 11)]
    
    # 4. 创建爬虫实例
    scraper = AsyncScraper(
        headers=headers,
        logger=logger,
        ssl_context=ssl_context,
        ua=ua
    )

    # 5. 启动爬虫
    await scraper.start()

    # 6. 爬取列表页
    logger.info("开始爬取列表页...")
    index_responses = await scraper.fetch_many(index_urls)

    # 7. 提取所有详情页URL
    all_next_urls = []
    for html in index_responses:
        if html:
            next_url = def_next_urls(html, logger)
            all_next_urls.extend(next_url)
    
    logger.info(f"共提取到 {len(all_next_urls)} 个详情页URL")

    # 8. 爬取详情页
    detail_urls = [url + detail for detail in all_next_urls]
    logger.info("开始爬取详情页...")
    detail_responses = await scraper.fetch_many(detail_urls)

    # 9. 关闭爬虫
    await scraper.close()

    # 10. 解析详情页数据
    all_data = []
    for response in detail_responses:
        if response:
            data = dict_data(response_text=response, logger=logger)
            if data:
                all_data.append(data)
        else:
            logger.error('详情页面解析失败')

    # 11. 保存数据
    with open('ssr2.json', 'w', encoding='utf-8') as f:
        json.dump(all_data, f, ensure_ascii=False, indent=2)
        logger.info(f'数据保存成功,共 {len(all_data)} 条')

    return '网站数据抓取完成,请查看保存的json文件'


if __name__ == '__main__':
    logger = logger_config()
    start = time.time()
    asyncio.run(main())
    end = time.time()
    logger.info(f'请求总用时: {end - start:.2f}秒')

代码设计

4.1 类型提示(Type Hints)

from typing import Optional, Dict, Tuple, List

def __init__(
    self,
    retry: int = 3,
    time_sleep: Tuple[float, float] = (0.5, 1.5),
    headers: Optional[Dict[str, str]] = None,
):
  • 提高代码可读性
  • IDE自动补全更智能

4.2 链接池管理

self.connector = aiohttp.TCPConnector(
    limit=limit,           # 全局最大连接数
    limit_per_host=limit_per_host,  # 单域名最大连接数
    ssl=self.ssl_context
)
  • 避免过多链接耗尽系统资源
  • 防止对目标服务器造成过大压力

4.3 信号并发控制

self.semaphore = asyncio.Semaphore(semaphore)

async with self.semaphore:
    # 受限的并发代码
  • 精准控制并发数量
  • 避免一次性创建过多任务

4.4 智能重试机制

# 服务器错误时重试,且等待时间递增
wait = random.uniform(*self.time_sleep) + attempt * 0.2
  • 指数退避策略
  • 随机抖动避免同时重试

五 运行结果

2024-01-01 10:00:00 - index.py - INFO - 爬虫引擎启动
2024-01-01 10:00:00 - index.py - INFO - 开始爬取列表页...
2024-01-01 10:00:03 - index.py - INFO - 共提取到 150 个详情页URL
2024-01-01 10:00:03 - index.py - INFO - 开始爬取详情页...
2024-01-01 10:00:15 - index.py - INFO - 成功解析电影: xxxxxxx
2024-01-01 10:00:15 - index.py - INFO - 成功解析电影:xxxxxxxxx
...
2024-01-01 10:00:30 - index.py - INFO - 数据保存成功,共 1502024-01-01 10:00:30 - index.py - INFO - 爬虫引擎已关闭
2024-01-01 10:00:30 - index.py - INFO - 请求总用时: 30.25

生成的 文件.json文件示例:

[
  {
    "name": "xxxxxxxxxx",
    "categories": 'xxxxxx',
    "categories_str": "xxxxxxx",
    "region": "xx",
    "duration": "142分钟",
    "score": "9.7",
    "drama": "一场谋杀案使银行家安迪蒙冤入狱..."
  },
  ...
]

六 后期扩展计划

  • 添加代理
  • 理解生产消费模型做请求队列
  • 数据持久化

七 总结

  • 代码解构清晰,每个类和方法职责明确
  • 配置灵活,初始化参数可调整
  • 资源管理规范,连接池,信号量,生命周期管理
  • 可复用性,类可以用于其他爬虫项目
  • 类型提示,提高代码可维护性
Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐