基于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文件里,启动时自动加载,不用硬编码到代码里。
三、单独启停与状态管理
要实现单独启停,得维护每个爬虫的运行状态:
- 用一个字典记录当前运行的爬虫进程(键为爬虫名,值为进程对象)
- 用轻量的SQLite数据库持久化爬虫的历史状态(运行中、已停止、已完成)、爬取数量、起止时间等
- 提供命令行或简单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()
额外优化建议
- 配置化管理:把70+爬虫的定时规则、运行参数都放到YAML/JSON配置文件里,避免硬编码,修改更方便
- API接口:如果需要远程控制,可以用Flask/FastAPI写个简单的接口,暴露
/run/<spider_name>、/stop/<spider_name>、/status等端点 - 资源监控:可以结合
psutil库监控每个爬虫进程的CPU、内存占用,补充到监控数据里 - 异常重试:给定时任务添加重试机制,比如爬虫失败后自动重试1-2次
内容的提问来源于stack exchange,提问作者NFB
相关产品推荐
相关产品推荐

