如何在Celery以worker模式启动时执行指定代码,避免客户端导入时触发
Celery仅在Worker启动时执行初始化代码的实现方案
原顶级作用域编写初始化代码的写法存在问题:只要对应模块被任意角色(包括发任务的客户端、Worker进程)导入,代码就会执行,不符合你的需求。你可以使用Celery官方提供的进程信号钩子实现要求:
方案说明
Celery提供了专门的Worker生命周期信号,只有对应Worker进程启动时才会触发信号绑定的函数,客户端导入任务模块时不会执行相关逻辑:
- 若你需要在Worker主进程仅初始化一次资源,使用
worker_init信号 - 若你使用prefork进程池,需要为每个工作子进程单独初始化数据库连接这类不可跨进程共享的资源,使用
worker_process_init信号(更适合你的SQLAlchemy引擎初始化场景,避免跨进程连接异常)
代码示例
from celery.signals import worker_process_init from celery import Celery from sqlalchemy import create_engine # 你的Celery应用初始化 celery_app = Celery(__name__) # 先声明资源变量,初始化前默认为None engine = None # 绑定信号触发的初始化函数 @worker_process_init.connect def init_worker_resources(sender, **kwargs): global engine # 仅Worker子进程启动时才会执行以下逻辑 engine = create_engine(str(POSTGRES_URL))
如果不想使用全局变量,也可以将资源绑定到Celery应用实例上:
@worker_process_init.connect def init_worker_resources(sender, **kwargs): # sender参数为当前Worker对应的Celery应用实例 sender.engine = create_engine(str(POSTGRES_URL)) # 任务中使用时直接从celery_app取即可 @celery_app.task def demo_task(): with celery_app.engine.connect() as conn: # 执行数据库操作 pass
内容的提问来源于stack exchange,提问作者The Fool
相关产品推荐
相关产品推荐

