基于FastAPI实现Selenium爬虫单线程任务队列的最佳实践咨询
Python实现生产者-消费者模式的最佳实践(适配Selenium爬虫API场景)
核心思路
你的Selenium爬虫需要独占激活窗口,单任务耗时3-4分钟,用任务队列解耦API请求与爬取执行是最优解:API作为生产者仅负责提交任务,后台消费者进程/线程逐个执行爬取,彻底避免窗口抢占导致的中断问题。
具体实现方案
1. 标准库queue轻量实现(单服务器原型场景)
无需额外依赖,适合快速验证思路:
- 生产者(API端点):用Flask/FastAPI编写接口,接收爬取请求后,将任务参数(如目标URL、爬取规则)存入
queue.Queue。 - 消费者(后台线程):启动单独立线程,循环从队列取任务,初始化Selenium实例执行爬取,完成后将结果存入数据库/缓存供API查询。
示例代码片段:
from queue import Queue import threading from selenium import webdriver from flask import Flask, request, jsonify app = Flask(__name__) task_queue = Queue() # 生产环境建议用Redis/数据库替代内存字典 task_results = {} def consumer_thread(): while True: task_id, params = task_queue.get() try: driver = webdriver.Chrome() driver.get(params["url"]) # 执行自定义爬取逻辑 crawl_data = {"page_title": driver.title} task_results[task_id] = {"status": "success", "data": crawl_data} except Exception as e: task_results[task_id] = {"status": "failed", "error": str(e)} finally: driver.quit() task_queue.task_done() # 启动后台消费者线程 threading.Thread(target=consumer_thread, daemon=True).start() @app.route("/submit-crawl", methods=["POST"]) def submit_task(): params = request.json task_id = f"task_{hash(frozenset(params.items()))}" task_results[task_id] = {"status": "pending"} task_queue.put((task_id, params)) return jsonify({"task_id": task_id, "msg": "任务已加入队列"}) @app.route("/get-result/<task_id>", methods=["GET"]) def fetch_result(task_id): return jsonify(task_results.get(task_id, {"status": "not_found"})) if __name__ == "__main__": app.run(host="0.0.0.0", port=5000)
2. Celery+Redis/RabbitMQ生产级实现
适合需要多节点扩展、任务重试、监控的场景:
- 生产者:API调用Celery的
delay()方法提交爬取任务。 - 消费者:启动Celery Worker进程,设置
--concurrency=1确保单进程单任务执行,避免窗口抢占。 - 结果存储:用Redis/RabbitMQ作为任务队列和结果后端,天然支持分布式场景。
示例配置与代码:
# celery_tasks.py from celery import Celery from selenium import webdriver # 初始化Celery,用Redis作为消息代理和结果后端 celery_app = Celery( "crawl_service", broker="redis://localhost:6379/0", backend="redis://localhost:6379/0" ) @celery_app.task(bind=True, retry_backoff=3) def execute_crawl(self, task_params): driver = None try: driver = webdriver.Chrome() driver.get(task_params["url"]) # 自定义爬取逻辑 return {"page_source_length": len(driver.page_source)} except Exception as e: # 异常自动重试,最多3次 self.retry(exc=e, max_retries=3) finally: if driver: driver.quit()
FastAPI端点示例:
from fastapi import FastAPI from celery_tasks import execute_crawl app = FastAPI() @app.post("/crawl") async def add_crawl_task(task_params: dict): task = execute_crawl.delay(task_params) return {"task_id": task.id, "status": "queued"} @app.get("/result/{task_id}") async def get_crawl_result(task_id: str): task = execute_crawl.AsyncResult(task_id) match task.state: case "PENDING": return {"status": "pending"} case "SUCCESS": return {"status": "success", "data": task.result} case _: return {"status": "failed", "error": str(task.info)}
关键注意事项
- 窗口独占保障:消费者必须单线程/单进程执行任务,Celery Worker需设置
--concurrency=1,标准库实现用单消费者线程,避免多窗口抢占激活状态。 - 资源清理:每个任务结束后必须关闭Selenium浏览器实例,防止内存泄漏。
- 状态持久化:生产环境禁止用内存字典存储任务结果,改用Redis、MySQL等持久化存储。
- 异常处理:添加任务重试机制,处理页面加载超时、元素定位失败等常见Selenium异常。
内容的提问来源于stack exchange,提问作者nextedoff
相关产品推荐
相关产品推荐

