欢迎光临
我们一直在努力

一个虫子的养成日记:python异步爬虫,基于任务丢列的生产者-消费者模式实战

Python异步爬虫高阶:基于任务队列的生产者-消费者模式实战

一、前言

在之前的文章中,我实现了基础异步爬虫和面向对象封装的爬虫,现在进一步升级,引入 任务队列(queue)和生产者 – 消费者模式 ,构建一个更加灵活,可扩展的异步爬虫框架。 项目特点

  • 基于·asyncio.queue 的任务调度
  • 多worker并发处理不同的任务
  • 异步下入文件(带锁)
  • 模块设计:爬虫核心,任务定义,worker,pipeling

二 . 为什么需要队列任务?

2.1 传统爬虫的局限

# 传统方式:先爬列表页,再爬详情页
list_responses = await fetch_many(list_urls)
detail_urls = extract(list_responses)
detail_responses = await fetch_many(detail_urls)

  • 串行执行:必须等待所有列表爬虫完成才能爬取详情页
  • 无法动态添加任务:url必须预先知道

2.2任务队列的优势

生产者(投放初始URL) → [任务队列] → 消费者Worker1 → 提取新URL → 再次投放

消费者Worker2 → 解析数据 → 消费者Worker3 → 保存

  • 解耦:任务生产者和消费分离
  • 动态性:可以随时添加任务到worker中
  • 负载:多个worker 并行处理
  • 可扩展: 轻松增加新的任务类型

完整代码

"""
File: index.py
慢速站点异步爬虫(任务驱动版)
作者:张志刚
"""

# 全局变量
BASE_URL = '这里不公开哈'

# 模块导入
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
import json

3.2 日志和ssl配置

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

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(object):
def __init__(self,
retry: int = 3,
time_sleep: Tuple[float, float] = (0.5, 1.5),
timeout=aiohttp.ClientTimeout(total=3),
headers: Dict[str, str] = None,
semaphore: Optional[int] = 5,
logger: Optional[logging.Logger] = None,
ssl_context: Optional[ssl.SSLContext] = None,
ua: Optional[UserAgent] = None,
limit_per_host: int = 10,
limit: int = 3,
):
"""
异步爬虫核心类

:param retry: 重试次数
:param time_sleep: 请求延时范围
:param timeout: 请求超时
:param headers: 请求头
:param semaphore: 并发控制数
:param logger: 日志实例
:param ssl_context: SSL配置
:param ua: 随机请求头生成器
:param limit_per_host: 单域名连接数
:param limit: 全局连接数
"""
self.retry = retry
self.time_sleep = time_sleep
self.timeout = timeout
self.headers = headers
self.ua = ua or UserAgent()
self.session = None
self.semaphore = asyncio.Semaphore(semaphore)
self.ssl_context = ssl_context or ssl.create_default_context()
self.logger = logger or logging.getLogger()

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

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:

# 成功响应
if 200 <= response.status < 300:
await asyncio.sleep(random.uniform(*self.time_sleep))
return await response.text()

# 需要重试的状态码(如429 Too Many Requests)
elif response.status < 429:
await asyncio.sleep(random.uniform(*self.time_sleep))
continue

except Exception as e:
self.logger.warning(f'异常: {e} -> 重试 {attempt + 1}')

return None

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.4 任务系统设计

这里说一下我的理解,这里就是一个工单,具体描述要告诉worker任务的类型,着重理解了一下这里,worker通过判断task_type 来完判断任务是做什么的~~~理解了这里估计就理解了worker了,我就是从这里理解了worker!!!

def make_task(task_type, data):
"""
创建任务(统一格式)

:param task_type: 任务类型('page'|'detail'|'response')
:param data: 任务数据
:return: 任务字典
"""
return {
'task_type': task_type,
'data': data,
}

类型作用生产者消费者
page 爬取列表页 主函数 worker解析出detail任务
detail 爬取详情页 worker worker解析出response任务
response 保存数据 worker worker调用pipeline保存

3.5 页面解析逻辑

这里不做过多描述,代码放在这里就可以了

def def_next_urls(response_text):
"""从列表页提取详情页URL"""
if not response_text:
return []

html = etree.HTML(response_text)
if html is None:
return []

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

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.6 异步 json 存储 pipeline

class JsonPipeline(object):
def __init__(self, file_name='result.json'):
"""
异步JSON存储管道

:param file_name: 输出文件名
"""
self.file_name = file_name
self.first = True
self.lock = asyncio.Lock() # 异步锁,防止多个Worker同时写入

# 初始化文件,写入JSON数组开始
with open(self.file_name, 'w', encoding='utf-8') as f:
f.write('[\\n')

async def save(self, item):
"""
异步保存单条数据

:param item: 要保存的数据字典
"""
async with self.lock: # 加锁保证写入安全
with open(self.file_name, 'a', encoding='utf-8') as f:
if not self.first:
f.write(',\\n')
json.dump(item, f, ensure_ascii=False, indent=4)
self.first = False

async def close(self):
"""关闭Pipeline,完成JSON数组"""
async with self.lock:
with open(self.file_name, 'a', encoding='utf-8') as f:
f.write('\\n]')

3.7 worker工作单元

这里说,save 做了一个单独的模块,而不是在 detail 这个模块中直接做保存,是为了,防止阻塞

async def worker(name, queue, scraper, logger, pipeline):
"""
消费者Worker:从队列获取任务并处理

:param name: Worker名称
:param queue: 任务队列
:param scraper: 爬虫实例
:param logger: 日志实例
:param pipeline: 存储管道
"""
while True:
# 1. 从队列获取任务
task = await queue.get()

# 2. 检查退出信号
if task is None:
logger.info(f'{name} 退出')
queue.task_done()
break

# 3. 解析任务
task_type = task['task_type']
url = task['data']

# ———— 类型1:列表页任务 —————
if task_type == 'page':
html = await scraper._fetch(url)

if html:
# 提取详情页URL
paths = def_next_urls(html)

# 生成新任务并放回队列
for p in paths:
full_url = BASE_URL + p
await queue.put(make_task('detail', full_url))

# ———— 类型2:详情页任务 —————
elif task_type == 'detail':
html = await scraper._fetch(url)

if html:
logger.info(f'详情页面抓取成功:{url},准备解析')
item = dict_data(response_text=html, logger=logger)
# 生成保存任务
await queue.put(make_task('response', item))

# ———— 类型3:数据保存任务 —————
elif task_type == 'response':
item = task['data']
logger.info(f'数据保存:{item["name"]}')
await pipeline.save(item)

# 4. 标记任务完成
queue.task_done()

3.8 主函数

async def main():
"""主函数:协调整个爬虫流程"""

# 1. 初始化配置
logger = logger_config()
ssl_context = ssl_config()
UA = UserAgent()

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',
}

# 2. 创建任务队列
queue = asyncio.Queue()

# 3. 创建爬虫实例
scraper = AsyncScraper(
headers=headers,
logger=logger,
ssl_context=ssl_context,
ua=UA
)

# 4. 创建Pipeline
pipeline = JsonPipeline()

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

# 6. 创建Worker协程(5个并发Worker)
workers = [
asyncio.create_task(
worker(f'worker-{i}', queue, scraper, logger, pipeline=pipeline)
)
for i in range(5)
]

# 7. 投放初始任务(10个列表页)
for i in range(1, 11):
url = BASE_URL + '/page/' + str(i)
await queue.put(make_task('page', url))

logger.info(f"初始任务投放完成: {queue.qsize()}")

# 8. 等待所有任务完成
await queue.join()

# 9. 发送退出信号给所有Worker
for _ in workers:
await queue.put(None)

# 10. 等待所有Worker退出
await asyncio.gather(*workers)

# 11. 关闭资源
await scraper.close()
await pipeline.close()

logger.info("爬虫结束")

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

四. 核心

4.1异步锁机制

self.lock = asyncio.Lock() # 异步锁

async with self.lock: # 保证同一时间只有一个Worker写入
# 文件写入操作

为什么需要锁?

  • 多个worker同时写入一个文件导致数据混乱
  • asyncio.lock()提供协程级别的互斥访问

5 版本对比

版本架构并发模式灵活性代码行数
v1基础版 过程式 batch_fetch ~150
v2基础版 类封装 Semaphore控制 ~200
v3基础版 生产者-消费者式 Queue + Workers底 ~280

6.2 任务队列优势

  • 动态任务生成:爬虫过程中发现新的url立即 put 进新任务
  • 负载均衡:空闲的worker自动获取新任务
  • 背压处理:队列满了自动阻塞生产者
  • 优雅退出:None信号控制worker退出
  • 7 完整代码

    """
    File: index.py
    慢速站点异步爬虫(任务驱动版)
    作者:张志刚
    """

    #################################################### 全局变量 ################################################

    BASE_URL = '不公开,怕被揍'

    ###################################### 模块 ##################################################################
    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
    import json
    ###################################### logging和ssl配置 #######################################################

    def logger_config():
    logging.basicConfig(level=logging.INFO,
    format='%(asctime)s – %(filename)s – %(levelname)s – %(message)s'
    )
    return logging.getLogger()

    def ssl_config():
    ssl_context = ssl.create_default_context()
    ssl_context.check_hostname = False
    ssl_context.verify_mode = ssl.CERT_NONE
    return ssl_context

    ###################################### 爬虫核心 #######################################################

    class AsyncScraper(object):
    def __init__(self,
    retry:int = 3,
    time_sleep:Tuple[float, float] = (0.5,1.5),
    timeout = aiohttp.ClientTimeout(total=3),
    headers:Dict[str, str]=None,
    semaphore:Optional[int]=5,
    logger:Optional[logging.Logger]=None,
    ssl_context:Optional[ssl.SSLContext]=None,
    ua:Optional[UserAgent]=None,
    limit_per_host:int = 10,
    limit:int=3,
    ):
    '''

    :param retry: 重试次数
    :param time_sleep: 请求延时
    :param timeout: 请求超时
    :param headers: 请求头
    :param semaphore: 并发线程
    :param logger:
    :param ssl_context:
    :param ua: 随机请求头
    :param limit_per_host:链接池
    :param limit: 单域链接
    '''
    self.retry = retry
    self.time_sleep = time_sleep
    self.timeout = timeout
    self.headers = headers
    self.ua = ua or UserAgent()
    self.session = None

    self.semaphore = asyncio.Semaphore(semaphore)

    self.ssl_context = ssl_context or ssl.create_default_context()
    self.logger = logger or logging.getLogger()

    self.connector = aiohttp.TCPConnector(ssl=self.ssl_context,
    limit_per_host=limit_per_host,
    limit=limit,
    )

    async def _fetch(self, url:str,):

    for attempt in range(self.retry):
    try:
    async with self.semaphore:
    curent_headers = self.headers.copy()
    curent_headers['User-Agent'] = self.ua.random

    self.logger.info(f'请求{url}')

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

    if 200 <= response.status < 300:
    await asyncio.sleep(random.uniform(*self.time_sleep))
    return await response.text()

    elif response.status < 429:
    await asyncio.sleep(random.uniform(*self.time_sleep))
    continue

    except Exception as e:

    self.logger.warning(f'异常: {e} -> 重试 {attempt + 1}')

    return None

    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('爬虫关闭')

    ############################################### 任务结构 #############################################

    def make_task(task_type,data):
    return {
    'task_type':task_type,
    'data':data,
    }

    ############################################### 解析 #############################################

    def def_next_urls(response_text):
    if not response_text:
    return []

    html = etree.HTML(response_text)
    if html is None:
    return []

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

    def dict_data(response_text,logger)>dict:

    '''
    详情页面主要信息提取
    :param response_text: 页面
    :param logger: logging.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}

    ######################################## 保存 json #################################################
    class JsonPipeline(object):
    def __init__(self,file_name='result.json'):

    self.file_name = file_name
    self.first = True
    self.lock = asyncio.Lock()

    with open(self.file_name,'w',encoding='utf-8') as f:

    f.write('[\\n')

    async def save(self,item):

    async with self.lock:
    with open(self.file_name,'a',encoding='utf-8') as f:
    if not self.first:
    f.write(',\\n')
    json.dump(item, f, ensure_ascii=False,indent=4)
    self.first = False

    async def close(self):

    async with self.lock:
    with open(self.file_name,'a',encoding='utf-8') as f:
    f.write('[\\n')

    ############################################### Worker #############################################

    async def worker(name,queue,scraper,logger,pipeline):

    while True:
    task = await queue.get()

    if task is None:
    logger.info(f'{name} 退出')
    queue.task_done()
    break

    task_type = task['task_type']
    url = task['data']

    # ———— PAGE —————

    if task_type == 'page':

    html = await scraper._fetch(url)

    if html:
    paths = def_next_urls(html)

    for p in paths:
    full_url = BASE_URL + p
    await queue.put(make_task('detail', full_url))

    # —————- DETAIL ————–

    elif task_type == 'detail':

    html = await scraper._fetch(url)

    if html:

    logger.info(f'详情页面抓取成功:{url},准备解析')
    item = dict_data(response_text=html, logger=logger)
    await queue.put(make_task('response', item))
    # ———————— save —————————————
    elif task_type == 'response':

    item = task['data']
    logger.info(f'数据保存:{item["name"]}')

    await pipeline.save(item)

    queue.task_done()

    # ——————————————————– 启动 ———————————————–
    async def main():

    logger = logger_config()
    ssl_context = ssl_config()
    UA = UserAgent()

    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',
    }

    queue = asyncio.Queue()

    scraper = AsyncScraper(
    headers=headers,
    logger=logger,
    ssl_context=ssl_context,
    ua=UA
    )

    pipeline = JsonPipeline()

    await scraper.start()

    logger.info("爬虫启动")

    workers = [
    asyncio.create_task(
    worker(f'worker-{i}', queue, scraper, logger, pipeline=pipeline)
    )
    for i in range(5)
    ]

    for i in range(1, 11):
    url = BASE_URL + '/page/' + str(i)
    await queue.put(make_task('page', url))

    logger.info(f"初始任务投放完成: {queue.qsize()}")

    await queue.join()

    for _ in workers:
    await queue.put(None)

    await asyncio.gather(*workers)

    # 8. 关闭资源
    await scraper.close()
    await pipeline.close()

    logger.info("爬虫结束")

    ###############################################
    if __name__ == '__main__':
    asyncio.run(main())

    张志刚 Happy coding! 2026-02-22

    赞(0)
    未经允许不得转载:171主机测评 » 一个虫子的养成日记:python异步爬虫,基于任务丢列的生产者-消费者模式实战
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址