FastAPI事件循环与Scrapy的Twisted线程是否存在兼容性问题?
问题背景
我有一个Scrapy爬虫,希望通过FastAPI的REST API端点触发。单独运行启动脚本完全正常,但集成到FastAPI接口后,爬虫仅首次触发成功,后续每次触发都会失败,日志中出现ValueError: signal only works in main thread of the main interpreter错误。
原始Scrapy启动脚本(单独运行正常)
from scrapy.crawler import CrawlerProcess from scrapy.utils.project import get_project_settings from tundra.spiders.spiderfool import SpiderfoolSpider def run_spiderfool(start_page=1, number_pages=2): process = CrawlerProcess(get_project_settings()) process.crawl(SpiderfoolSpider, start_page=start_page, number_pages=number_pages) process.start() if __name__ == "__main__": run_spiderfool(1, 2)
首次集成FastAPI的代码(仅首次触发成功)
@app.get("/api/trigger/spiderfool") def trigger_spider( start_page: int = Query(1, ge=1), number_pages: int = Query(2, ge=1) ): try: logger.info(f"Starting spider with start_page={start_page}, number_pages={number_pages}") run_spiderfool(start_page=start_page, number_pages=number_pages) return {"message": f"Success! Triggered spider to crawl Motley Fool's transcripts collection from page {start_page} for {number_pages} pages!"} except Exception as e: logger.error(f"Failure! Triggered no spider! Err: {str(e)}") return {"error": "Failure! Triggered no spider!", "details": str(e)}
尝试过的无效方案
- 安装Twisted的
asyncioreactor.install(),将端点改为异步方法并通过asyncio.create_task调用 - 使用
nest_asyncio修补事件循环
当前可行但有缺陷的方案(subprocess)
通过subprocess.Popen启动Scrapy命令,解决了重复触发问题,但无法跟踪爬虫完成状态:
import subprocess def run_spiderfool(start_page=1, number_pages=2): subprocess.Popen(['scrapy', 'crawl', 'spiderfool', f'-a', f'start_page={start_page}', f'-a', f'number_pages={number_pages}'])
核心原因:FastAPI事件循环与Scrapy/Twisted的信号机制冲突
FastAPI基于uvloop(默认)或标准asyncio事件循环运行,而Scrapy依赖的Twisted框架在启动CrawlerProcess.start()时,会尝试注册只能在主线程生效的信号处理器(比如SIGINT、SIGTERM)。首次调用后,信号处理器已经被注册,且Twisted的reactor无法被重复启动/重置,后续调用时就会触发signal only works in main thread的错误,导致爬虫无法再次运行。
你的异步改造方案无效,是因为run_in_executor只是把爬虫运行在另一个线程,但Twisted的reactor一旦启动就无法在同一个进程中再次初始化,本质问题没解决。
优化解决方案
1. 改进subprocess方案,增加状态跟踪
可以通过记录进程PID、输出日志,或结合文件/数据库标记任务状态,弥补无法跟踪的缺陷:
import subprocess import os from datetime import datetime # 内存存储任务信息,生产环境建议用Redis或数据库 spider_tasks = {} def run_spiderfool(start_page=1, number_pages=2): # 生成唯一任务ID task_id = f"spiderfool_{datetime.now().strftime('%Y%m%d_%H%M%S')}" # 日志文件路径 log_path = f"logs/{task_id}.log" # 启动进程并保存PID proc = subprocess.Popen( ['scrapy', 'crawl', 'spiderfool', '-a', f'start_page={start_page}', '-a', f'number_pages={number_pages}'], stdout=open(log_path, 'w'), stderr=subprocess.STDOUT ) # 存储任务信息 spider_tasks[task_id] = {"pid": proc.pid, "log_path": log_path, "status": "running"} return task_id # 新增查询状态的端点 @app.get("/api/spider/status/{task_id}") def get_spider_status(task_id: str): if task_id not in spider_tasks: return {"error": "Task not found"} task_info = spider_tasks[task_id] try: os.kill(task_info["pid"], 0) # 检查进程是否存活 return {"task_id": task_id, "status": "running", "log_path": task_info["log_path"]} except OSError: task_info["status"] = "completed" return {"task_id": task_id, "status": "completed", "log_path": task_info["log_path"]}
2. 使用任务队列(如Celery)解耦FastAPI与Scrapy
这是更健壮的生产级方案,彻底解决进程冲突问题,同时原生支持任务状态跟踪:
Celery配置(tasks.py)
from celery import Celery from scrapy.crawler import CrawlerProcess from scrapy.utils.project import get_project_settings from tundra.spiders.spiderfool import SpiderfoolSpider # 用Redis作为消息代理,可替换为其他支持的代理 app = Celery('spider_tasks', broker='redis://localhost:6379/0') @app.task(bind=True) def crawl_spiderfool(self, start_page=1, number_pages=2): process = CrawlerProcess(get_project_settings()) process.crawl(SpiderfoolSpider, start_page=start_page, number_pages=number_pages) process.start() return "Crawl completed successfully"
FastAPI端点
from tasks import crawl_spiderfool @app.get("/api/trigger/spiderfool") def trigger_spider( start_page: int = Query(1, ge=1), number_pages: int = Query(2, ge=1) ): task = crawl_spiderfool.delay(start_page, number_pages) return {"task_id": task.id, "message": "Crawl task started"} @app.get("/api/spider/status/{task_id}") def get_task_status(task_id: str): task = crawl_spiderfool.AsyncResult(task_id) if task.state == 'PENDING': response = {"status": "Pending..."} elif task.state == 'SUCCESS': response = {"status": "Completed", "result": task.result} else: response = {"status": task.state, "error": str(task.info)} return response
3. 自定义Scrapy reactor管理(不推荐生产用)
如果坚持在同一个进程内运行,可以尝试手动重置Twisted reactor,但这种方式容易出现线程安全问题,仅适合测试场景:
from scrapy.crawler import CrawlerProcess from scrapy.utils.project import get_project_settings from twisted.internet import reactor def run_spiderfool(start_page=1, number_pages=2): process = CrawlerProcess(get_project_settings()) process.crawl(SpiderfoolSpider, start_page=start_page, number_pages=number_pages) # 手动启动reactor,运行完成后停止 def callback(_): reactor.stop() process.start(False) # 不阻塞主线程 reactor.addSystemEventTrigger('after', 'shutdown', callback) reactor.run()
内容的提问来源于stack exchange,提问作者Sun Bee

