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

基于CrawlProcess的Scrapy并发爬虫排队与监控方案咨询

嘿,针对你这个要管理70多个活跃开发爬虫、又不想用Scrapyd的需求,我整理了一套基于CrawlProcess的落地方案,完全贴合你的定时、并发、单独启停和监控要求,咱们一步步来拆解:

一、核心架构:CrawlProcess+多进程实现隔离与并发

Scrapy的CrawlProcess本身是单爬虫的运行容器,要同时跑多个爬虫且支持单独启停,最好用多进程来隔离每个爬虫的运行环境——这样单个爬虫崩溃不会影响其他,也方便单独终止某个进程。

核心思路:

  • 封装一个SpiderManager类,统一管理所有爬虫的生命周期
  • 用multiprocessing.Process给每个爬虫单独启动一个CrawlProcess
  • 用信号量控制同时运行的爬虫数不超过4个
二、定时任务:APScheduler实现每日/每周调度

APScheduler是轻量且灵活的定时任务框架,完全不需要复杂部署,代码修改后直接重启就能生效,特别适合活跃开发的项目。

你可以用Cron表达式配置定时规则,比如:

  • 每日凌晨2点运行某爬虫:hour=2
  • 每周日凌晨3点运行某爬虫:day_of_week='sun', hour=3

为了方便管理70+爬虫的定时规则,建议把配置写在YAML/JSON文件里,启动时自动加载,不用硬编码到代码里。

三、单独启停与状态管理

要实现单独启停,得维护每个爬虫的运行状态:

  1. 用一个字典记录当前运行的爬虫进程(键为爬虫名,值为进程对象)
  2. 用轻量的SQLite数据库持久化爬虫的历史状态(运行中、已停止、已完成)、爬取数量、起止时间等
  3. 提供命令行或简单API接口,接收run <爬虫名>/stop <爬虫名>指令,调用对应的方法
四、爬虫监控:信号钩子+日志双维度监控

监控每个爬虫的运行情况,可以从两个方面入手:

  • 状态监控:通过Scrapy的信号钩子(spider_opened/spider_closed/spider_error)收集爬虫的起止时间、爬取条目数、错误信息,存入SQLite,随时可以查询
  • 日志监控:在Scrapy的settings里配置每个爬虫的日志单独输出到文件(比如LOG_FILE = f'logs/{spider_name}_{datetime.now().strftime("%Y%m%d")}.log'),方便排查问题;还可以给logging添加告警handler,当爬虫出现ERROR级别的日志时自动发邮件提醒
完整代码示例

下面是一个简化版的实现,你可以根据自己的项目需求扩展:

import scrapy
from scrapy.crawler import CrawlerProcess
from scrapy.utils.project import get_project_settings
from apscheduler.schedulers.background import BackgroundScheduler
import multiprocessing
import sqlite3
from datetime import datetime
import logging

# 配置基础日志
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)

class SpiderManager:
    def __init__(self, max_concurrent=4):
        self.settings = get_project_settings()
        # 动态配置每个爬虫的单独日志路径
        self.settings['LOG_FILE'] = lambda spider: f'logs/{spider.name}_{datetime.now().strftime("%Y%m%d")}.log'
        self.max_concurrent = max_concurrent
        self.scheduler = BackgroundScheduler(timezone='Asia/Shanghai')
        self.running_spiders = {}  # 记录运行中的爬虫进程
        self.status_db = sqlite3.connect('spider_status.db', check_same_thread=False)
        self._init_status_db()

    def _init_status_db(self):
        """初始化状态数据库表"""
        cursor = self.status_db.cursor()
        cursor.execute('''CREATE TABLE IF NOT EXISTS spider_status
                          (spider_name TEXT PRIMARY KEY, 
                           status TEXT DEFAULT 'idle',
                           start_time TEXT,
                           end_time TEXT,
                           item_count INTEGER DEFAULT 0,
                           error_msg TEXT)''')
        self.status_db.commit()

    def _update_status(self, spider_name, **kwargs):
        """更新爬虫状态到数据库"""
        set_clause = ', '.join([f"{k} = ?" for k in kwargs.keys()])
        values = list(kwargs.values()) + [spider_name]
        cursor = self.status_db.cursor()
        cursor.execute(f'''UPDATE spider_status SET {set_clause} WHERE spider_name = ?''', values)
        if cursor.rowcount == 0:
            # 新爬虫插入记录
            cols = ', '.join(kwargs.keys())
            placeholders = ', '.join(['?' for _ in kwargs.values()])
            cursor.execute(f'''INSERT INTO spider_status ({cols}, spider_name) VALUES ({placeholders}, ?)''', values)
        self.status_db.commit()

    def _spider_opened(self, spider):
        """爬虫启动时的钩子函数"""
        self._update_status(spider.name, status='running', start_time=str(datetime.now()))
        self.running_spiders[spider.name] = multiprocessing.current_process()
        logger.info(f"爬虫 {spider.name} 已启动")

    def _spider_closed(self, spider, reason):
        """爬虫关闭时的钩子函数"""
        item_count = spider.crawler.stats.get_value('item_scraped_count', 0)
        status = 'finished' if reason == 'finished' else 'stopped'
        self._update_status(spider.name, status=status, end_time=str(datetime.now()), item_count=item_count)
        if spider.name in self.running_spiders:
            del self.running_spiders[spider.name]
        logger.info(f"爬虫 {spider.name} 已关闭,原因:{reason},爬取条目数:{item_count}")

    def _spider_error(self, failure, response, spider):
        """爬虫出错时的钩子函数"""
        error_msg = str(failure.value)
        self._update_status(spider.name, error_msg=error_msg)
        logger.error(f"爬虫 {spider.name} 出错:{error_msg}")

    def run_spider(self, spider_name):
        """启动单个爬虫"""
        if len(self.running_spiders) >= self.max_concurrent:
            logger.warning(f"已达到最大并发数 {self.max_concurrent},{spider_name} 进入等待队列")
            return

        # 从spiders模块加载爬虫类(根据你的项目结构调整)
        spider_cls = None
        try:
            from spiders import __all__ as spider_list
            for mod_name in spider_list:
                if mod_name.lower() == spider_name.lower():
                    module = __import__(f'spiders.{mod_name}', fromlist=[''])
                    spider_cls = getattr(module, mod_name)
                    break
        except Exception as e:
            logger.error(f"加载爬虫 {spider_name} 失败:{str(e)}")
            return

        if not spider_cls:
            logger.error(f"未找到爬虫 {spider_name}")
            return

        # 创建CrawlProcess并绑定信号钩子
        process = CrawlerProcess(self.settings)
        process.signals.connect(self._spider_opened, signal=scrapy.signals.spider_opened)
        process.signals.connect(self._spider_closed, signal=scrapy.signals.spider_closed)
        process.signals.connect(self._spider_error, signal=scrapy.signals.spider_error)

        # 用多进程运行,避免阻塞主进程
        p = multiprocessing.Process(target=process.start, name=f"Spider-{spider_name}")
        p.start()
        self.running_spiders[spider_name] = p

    def stop_spider(self, spider_name):
        """停止单个爬虫"""
        if spider_name not in self.running_spiders:
            logger.warning(f"爬虫 {spider_name} 未在运行中")
            return

        p = self.running_spiders[spider_name]
        p.terminate()
        p.join(timeout=10)
        if p.is_alive():
            p.kill()
        self._update_status(spider_name, status='stopped', end_time=str(datetime.now()))
        del self.running_spiders[spider_name]
        logger.info(f"爬虫 {spider_name} 已强制停止")

    def load_schedule_from_config(self, config_path):
        """从YAML配置文件加载定时任务"""
        import yaml
        with open(config_path, 'r', encoding='utf-8') as f:
            schedule_config = yaml.safe_load(f)
        for spider_name, cron_rule in schedule_config.items():
            self.scheduler.add_job(
                self.run_spider,
                'cron',
                args=[spider_name],
                **cron_rule,
                id=f"job-{spider_name}"
            )
        logger.info(f"已加载 {len(schedule_config)} 个定时任务")

    def start(self):
        """启动管理器主程序"""
        self.scheduler.start()
        logger.info("爬虫管理器已启动,等待指令或定时任务触发")
        # 命令行交互接口
        try:
            while True:
                cmd = input("\n输入指令(run <爬虫名>/stop <爬虫名>/status/exit):").strip()
                if not cmd:
                    continue
                parts = cmd.split(maxsplit=1)
                if parts[0] == 'run':
                    if len(parts) < 2:
                        print("请指定爬虫名,示例:run SpiderA")
                        continue
                    self.run_spider(parts[1])
                elif parts[0] == 'stop':
                    if len(parts) < 2:
                        print("请指定爬虫名,示例:stop SpiderA")
                        continue
                    self.stop_spider(parts[1])
                elif parts[0] == 'status':
                    cursor = self.status_db.cursor()
                    cursor.execute('SELECT spider_name, status, start_time, end_time, item_count FROM spider_status')
                    rows = cursor.fetchall()
                    print("\n===== 爬虫状态列表 =====")
                    for row in rows:
                        print(f"名称:{row[0]} | 状态:{row[1]} | 开始时间:{row[2] or '-'} | 结束时间:{row[3] or '-'} | 爬取数:{row[4]}")
                elif parts[0] == 'exit':
                    logger.info("正在关闭所有爬虫和调度器...")
                    for spider_name in list(self.running_spiders.keys()):
                        self.stop_spider(spider_name)
                    self.scheduler.shutdown()
                    self.status_db.close()
                    print("已退出")
                    break
                else:
                    print("未知指令,可用指令:run/stop/status/exit")
        except KeyboardInterrupt:
            self.start()

if __name__ == '__main__':
    manager = SpiderManager(max_concurrent=4)
    # 示例:从config.yaml加载定时任务,配置文件内容参考:
    # SpiderA:
    #   hour: 2
    # SpiderB:
    #   day_of_week: 'sun'
    #   hour: 3
    # manager.load_schedule_from_config('config.yaml')
    manager.start()
额外优化建议
  1. 配置化管理:把70+爬虫的定时规则、运行参数都放到YAML/JSON配置文件里,避免硬编码,修改更方便
  2. API接口:如果需要远程控制,可以用Flask/FastAPI写个简单的接口,暴露/run/<spider_name>、/stop/<spider_name>、/status等端点
  3. 资源监控:可以结合psutil库监控每个爬虫进程的CPU、内存占用,补充到监控数据里
  4. 异常重试:给定时任务添加重试机制,比如爬虫失败后自动重试1-2次

内容的提问来源于stack exchange,提问作者NFB

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 03:42:55