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

FastAPI WebSocket结合Celery优化多用户性能问题求助

基于Celery的性能优化方案

当前问题分析

你的代码通过每个WebSocket连接启动独立进程爬取400个站点,多用户访问时会瞬间产生大量进程,抢占CPU、网络资源,导致整体性能骤降。用Celery可以把任务处理和Web服务解耦,通过Worker池控制并发,避免资源耗尽。

具体优化步骤

1. 配置Celery实例

先在项目中初始化Celery,用Redis作为消息中间件(轻量易部署):

# app/celery_config.py
from celery import Celery

celery_app = Celery(
    'sherlock_tasks',
    broker='redis://localhost:6379/0',
    backend='redis://localhost:6379/0'
)

# 配置Worker参数,控制并发与任务生命周期
celery_app.conf.update(
    worker_prefetch_multiplier=1,
    worker_max_tasks_per_child=100,
)

2. 重构爬取任务为Celery任务

把原来的main函数改造成Celery任务,执行过程中将中间结果发布到Redis指定频道(用任务ID做频道名,避免用户消息串流):

# app/sherlock.py
import redis
import json
from app.celery_config import celery_app

# 初始化Redis客户端用于发布中间结果
redis_client = redis.Redis(host='localhost', port=6379, db=0)

@celery_app.task(bind=True)
def check_username_task(self, username):
    # 保留原有的站点爬取逻辑,替换进程队列逻辑为Redis发布
    for site_result in crawl_sites(username):
        # 把结果序列化后发布到对应频道
        redis_client.publish(f"task_{self.request.id}", json.dumps(site_result))
    
    # 任务结束时发送完成标记
    redis_client.publish(f"task_{self.request.id}", json.dumps({"status": "completed"}))
    return "done"

# 原有的站点爬取函数(示例)
def crawl_sites(username):
    # 这里是遍历400个站点的逻辑,逐个返回结果
    pass

3. 修改WebSocket路由逻辑

WebSocket不再直接启动进程,而是提交Celery任务,订阅Redis频道获取实时结果并推送给前端:

# app/routers/websocket.py
import asyncio
import json
from fastapi import APIRouter, WebSocket, WebSocketDisconnect
from app.celery_config import celery_app
from app.sherlock import check_username_task

router = APIRouter()

# 异步Redis客户端用于订阅频道
async_redis = redis.asyncio.Redis(host='localhost', port=6379, db=0)

@router.websocket_route("/{username}")
async def websocket_endpoint(websocket: WebSocket):
    await websocket.accept()
    username = websocket.path_params["username"]
    
    # 提交Celery任务
    task = check_username_task.delay(username)
    task_id = task.id
    channel_name = f"task_{task_id}"
    
    # 订阅Redis频道
    pubsub = async_redis.pubsub()
    await pubsub.subscribe(channel_name)
    
    try:
        while True:
            # 异步等待频道消息
            message = await pubsub.get_message(ignore_subscribe_messages=True, timeout=1)
            if message:
                result_dict = json.loads(message['data'])
                
                if result_dict.get("status") == "completed":
                    await websocket.send_json(result_dict)
                    await websocket.close()
                    break
                await websocket.send_json(result_dict)
            
            # 检查任务是否异常终止
            if task.state == 'FAILURE':
                await websocket.send_json({"status": "error", "message": "任务执行失败"})
                await websocket.close()
                break
    except WebSocketDisconnect:
        # 用户断开连接时,终止对应的Celery任务
        celery_app.control.revoke(task_id, terminate=True)
        await pubsub.unsubscribe(channel_name)
    finally:
        await pubsub.close()

4. 启动Celery Worker

在项目根目录运行命令启动Worker,根据服务器CPU核心数设置并发数:

# 例如4核服务器设置并发数为4
celery -A app.celery_config worker --loglevel=info --concurrency=4

优化效果说明

  • 原来每个用户请求占用一个进程,现在所有任务由Celery Worker池统一处理,并发数可精准控制,不会因用户量激增耗尽系统资源。
  • Redis Pub/Sub实现了任务结果的实时推送,效果和原进程队列一致,但支持分布式部署。
  • 任务可持久化,即使Web服务重启,未完成的任务依然会在Worker中继续执行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 08:01:32