如何将Celery与dependency_injector结合?FastAPI依赖复用问题
问题解决与优化方案
错误原因
你遇到的ImportError是因为celery.py中的container定义在if __name__ == "__main__"代码块内,只有直接运行celery.py时才会执行该块。但Celery启动worker时,会导入tasks.py,此时celery.py作为模块被加载,__name__不等于__main__,导致container未定义,无法导入。
高效复用依赖的解决方案
我们可以利用Celery的worker_init信号,让每个worker进程启动时仅初始化一次容器,所有该进程处理的任务共享这个容器实例,避免重复初始化带来的性能损耗。
1. 修改celery.py
from celery import Celery from celery.signals import worker_init from app.core.config import get_app_settings from app.core.containers import Container settings = get_app_settings() worker = Celery( "worker", backend=settings.celery_backend_url, broker=settings.celery_broker_url, include=["app.worker.tasks"], ) worker.conf.update( result_expires=3600, ) @worker_init.connect def init_container(sender, **kwargs): # 每个worker进程启动时初始化一次容器 container = Container() container.config.from_pydantic(settings) # 将容器挂载到worker实例上,供任务调用 sender.container = container if __name__ == "__main__": worker.start()
2. 修改tasks.py
import asyncio from celery import Task import app.db.postgres.repositories as pg_repositories from .celery import worker class Task1(Task): def __init__(self, *args, **kwargs): super().__init__(*args, **kwargs) # 从worker实例获取已初始化的容器,复用依赖 self.task_repository: pg_repositories.TaskRepository = self.app.container.task_repository() async def task1(self): # 任务业务逻辑示例 # tasks = await self.task_repository.get_all() ... @worker.task(ignore_result=False, bind=True, base=Task1, name="tasks.task1") def task1(self): asyncio.run(self.task1())
方案优势
- 避免重复初始化:每个worker进程仅初始化一次容器,所有该进程内的任务共享容器中的依赖实例(如数据库连接池、仓库对象),大幅提升效率。
- 进程安全:Celery采用多进程模型,每个worker进程拥有独立的容器实例,避免跨进程的资源冲突。
- 符合Celery设计:通过
worker_init信号和self.app访问容器,逻辑清晰且贴合Celery的运行机制。
内容的提问来源于stack exchange,提问作者TimeTREE
相关产品推荐
相关产品推荐

