如何在FastAPI与Gunicorn中实现多API请求排队处理?
在FastAPI中实现API请求排队(解决Selenium多请求冲突问题)
Selenium的WebDriver实例并非线程安全,多请求并发时直接创建或共享实例会导致资源冲突、报错,同时耗时任务的并发处理也会耗尽系统资源。以下是几种生产可用的解决方案:
方案1:Celery异步任务队列(生产级推荐)
通过Celery将耗时的Selenium任务托管到后台队列,请求仅提交任务并获取任务ID,后续通过ID查询处理结果,彻底避免并发冲突。
步骤&代码示例
- 安装依赖:
pip install celery redis fastapi uvicorn
- 配置Celery任务(
celery_app.py):
from celery import Celery from selenium import webdriver from selenium.webdriver.chrome.options import Options # 初始化Celery,用Redis做消息代理和结果存储 celery = Celery( 'selenium_tasks', broker='redis://localhost:6379/0', backend='redis://localhost:6379/0' ) def init_chrome_driver(): # 配置无头Chrome,避免界面开销 chrome_options = Options() chrome_options.add_argument('--headless=new') chrome_options.add_argument('--no-sandbox') chrome_options.add_argument('--disable-dev-shm-usage') return webdriver.Chrome(options=chrome_options) @celery.task(bind=True) def run_selenium_task(self, target_url): driver = init_chrome_driver() try: driver.get(target_url) page_title = driver.title return {"status": "success", "result": page_title} except Exception as e: return {"status": "failed", "error": str(e)} finally: # 确保任务完成后关闭浏览器,释放资源 driver.quit()
- 编写FastAPI接口(
main.py):
from fastapi import FastAPI from celery_app import celery app = FastAPI() @app.post("/submit-scrape") async def submit_scrape_task(url: str): # 提交任务到Celery队列 task = run_selenium_task.delay(url) return {"task_id": task.id} @app.get("/task-result/{task_id}") async def get_task_result(task_id: str): task = celery.AsyncResult(task_id) match task.state: case 'PENDING': return {"status": "任务等待处理"} case 'SUCCESS': return {"status": "任务完成", "data": task.result} case _: return {"status": "任务失败", "error": str(task.info)}
- 启动服务:
- 先启动Redis服务
- 启动Celery Worker:
celery -A celery_app worker --loglevel=info - 启动FastAPI:
uvicorn main:app --reload
方案2:本地内存队列+后台线程(小型场景适用)
如果不想引入外部中间件,可使用Python内置队列配合后台线程实现串行处理,适合低流量场景。
代码示例
from fastapi import FastAPI, BackgroundTasks from queue import Queue from selenium import webdriver from selenium.webdriver.chrome.options import Options import threading import uuid import time app = FastAPI() task_queue = Queue() task_results = {} def queue_worker(): while True: task_id, url = task_queue.get() try: chrome_options = Options() chrome_options.add_argument('--headless=new') chrome_options.add_argument('--no-sandbox') chrome_options.add_argument('--disable-dev-shm-usage') driver = webdriver.Chrome(options=chrome_options) driver.get(url) task_results[task_id] = {"status": "success", "result": driver.title} except Exception as e: task_results[task_id] = {"status": "failed", "error": str(e)} finally: driver.quit() task_queue.task_done() time.sleep(0.5) # 启动后台工作线程 threading.Thread(target=queue_worker, daemon=True).start() @app.post("/submit-task") async def submit_task(url: str): task_id = str(uuid.uuid4()) task_queue.put((task_id, url)) task_results[task_id] = {"status": "pending"} return {"task_id": task_id} @app.get("/task-status/{task_id}") async def get_task_status(task_id: str): return task_results.get(task_id, {"status": "未知任务ID"})
注意:该方案仅适用于单进程部署,多Gunicorn Worker会各自维护独立队列,无法统一调度。
方案3:WebDriver实例池(复用资源)
维护一个固定数量的WebDriver实例池,控制并发数,避免频繁创建销毁实例带来的开销,同时限制并发请求数。
代码示例
from fastapi import FastAPI from selenium import webdriver from selenium.webdriver.chrome.options import Options from concurrent.futures import ThreadPoolExecutor from functools import lru_cache import uuid import asyncio app = FastAPI() # 维护最多5个WebDriver实例 @lru_cache(maxsize=5) def get_webdriver(): chrome_options = Options() chrome_options.add_argument('--headless=new') chrome_options.add_argument('--no-sandbox') chrome_options.add_argument('--disable-dev-shm-usage') return webdriver.Chrome(options=chrome_options) # 限制最多5个并发任务 executor = ThreadPoolExecutor(max_workers=5) task_cache = {} @app.post("/scrape") async def scrape_url(url: str): task_id = str(uuid.uuid4()) def run_task(): try: driver = get_webdriver() driver.get(url) return {"status": "success", "result": driver.title} except Exception as e: return {"status": "failed", "error": str(e)} future = executor.submit(run_task) task_cache[task_id] = future return {"task_id": task_id, "status": "processing"} @app.get("/scrape-result/{task_id}") async def get_scrape_result(task_id: str): future = task_cache.get(task_id) if not future: return {"status": "unknown task"} if future.done(): return future.result() else: return {"status": "processing"}
关键注意事项
- Gunicorn多Worker部署时,必须使用共享队列(如Redis)统一管理任务,否则每个Worker会独立处理请求,仍可能引发Selenium冲突。
- 务必开启Selenium的无头模式,避免弹出浏览器窗口,降低系统资源占用。
- 所有方案中都要确保WebDriver实例在任务完成后被关闭,防止内存泄漏。
内容的提问来源于stack exchange,提问作者Saurabh Nakoti
相关产品推荐
相关产品推荐

