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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 05:22:04