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
相关产品推荐
相关产品推荐

