如何在FastAPI与Celery Beat间共享ConnectionManager类实例?
解决FastAPI+Celery无法共享ConnectionManager实例的问题
核心原因:Celery Worker 和 FastAPI 应用是完全独立的进程(甚至在Docker部署中是不同容器),它们的内存空间相互隔离。你在FastAPI进程中创建的manager实例,在Celery进程里根本无法访问——Celery执行任务时会初始化一个全新的ConnectionManager实例,自然看到空字典。
方案1:用共享存储(如Redis)存储连接状态
把进程内的字典替换为Redis这类共享存储,让FastAPI和Celery都能读写同一数据源。这是最稳妥、可扩展的方案,尤其适合多容器/多实例部署。
修改ConnectionManager类
import redis import time from fastapi import WebSocket class ConnectionManager: def __init__(self): # 连接到Redis容器(假设Redis服务名是redis) self.redis = redis.Redis(host="redis", port=6379, db=0, decode_responses=True) self.connections_key = "active_websockets" async def connect(self, websocket: WebSocket, client_id: str): await websocket.accept() # 存储客户端ID和最后活跃时间戳 self.redis.hset(self.connections_key, client_id, str(int(time.time()))) def disconnect(self, client_id: str): # 移除断开的连接 self.redis.hdel(self.connections_key, client_id) def clean_stale_connections(self, timeout_hours: int = 1): timeout_seconds = timeout_hours * 3600 current_time = int(time.time()) cleaned_count = 0 # 遍历所有连接,清理超时的 all_connections = self.redis.hgetall(self.connections_key) for client_id, last_active_str in all_connections.items(): last_active = int(last_active_str) if current_time - last_active > timeout_seconds: self.redis.hdel(self.connections_key, client_id) cleaned_count += 1 return cleaned_count
Celery定时任务
from celery import Celery from your_app.module import ConnectionManager # 配置Celery使用Redis作为 broker celery = Celery('cleanup_tasks', broker='redis://redis:6379/0') @celery.task def clean_stale_websockets(): manager = ConnectionManager() return manager.clean_stale_connections()
方案2:让Celery调用FastAPI内部接口执行清理
如果不想引入外部存储,可以让Celery通过HTTP请求触发FastAPI进程内的清理逻辑——这样清理代码依然在FastAPI进程中执行,能直接访问到原有的manager实例。
给FastAPI添加内部清理接口
from fastapi import FastAPI, WebSocket, Depends, HTTPException from fastapi.security import APIKeyHeader import os app = FastAPI() manager = ConnectionManager() # 从环境变量读取内部API密钥,避免硬编码 INTERNAL_API_KEY = os.getenv("INTERNAL_API_KEY", "your-secure-secret") api_key_header = APIKeyHeader(name="X-Internal-API-Key", auto_error=False) def validate_internal_access(api_key: str = Depends(api_key_header)): if api_key != INTERNAL_API_KEY: raise HTTPException(status_code=403, detail="Unauthorized") return api_key # 仅允许内部调用的清理接口 @app.post("/internal/clean-connections", dependencies=[Depends(validate_internal_access)]) async def run_connection_cleanup(): cleaned = manager.clean_stale_connections() return {"cleaned_connections": cleaned}
Celery任务改为发送HTTP请求
import requests from celery import Celery import os celery = Celery('cleanup_tasks', broker='redis://redis:6379/0') FASTAPI_URL = os.getenv("FASTAPI_URL", "http://fastapi:8000/internal/clean-connections") INTERNAL_API_KEY = os.getenv("INTERNAL_API_KEY", "your-secure-secret") @celery.task def trigger_cleanup(): headers = {"X-Internal-API-Key": INTERNAL_API_KEY} try: response = requests.post(FASTAPI_URL, headers=headers) response.raise_for_status() return response.json() except requests.exceptions.RequestException as e: return {"error": str(e)}
为什么之前的方法无效?
- 用
@task或@shared_task只是注册任务,但Celery Worker运行时会启动独立进程,无法直接访问FastAPI进程的内存。 - 全局变量仅存在于单个进程的内存空间中,跨进程完全不共享,所以Celery里的全局
manager是全新的空实例。
内容的提问来源于stack exchange,提问作者Daniel Nochess
相关产品推荐
相关产品推荐

