通过RabbitMQ监听执行Celery任务时遭遇循环依赖问题求助
解决Celery循环依赖及相关问题
核心问题:模块命名冲突
错误根源是你将自己的Celery配置文件命名为celery.py,与Celery官方库的模块名重名。当class_task.py执行from celery import app时,Python会优先导入当前目录下的celery.py,而该文件又试图从官方Celery库导入Celery类,形成循环导入,最终触发报错。
解决步骤
1. 重命名Celery配置文件
将celery.py改为celery_app.py,避免与官方库冲突,修改后的文件内容:
import os from celery import Celery os.environ.setdefault("DJANGO_SETTINGS_MODULE", "projectTest.settings") app = Celery( "appTest", broker=os.environ.get("RABBITMQ_CONN_STRING"), backend="rpc://", ) # 可选配置,参考Celery官方文档 app.conf.update( result_expires=3600, ) app.config_from_object("appTest.celeryConfig", namespace="CELERY") app.autodiscover_tasks() if __name__ == "__main__": app.start()
2. 修正class_task.py的导入问题
- 需从自定义配置文件导入Celery实例
app,而非官方celery库 - 补充缺失的
get_task_logger导入语句
修改后的class_task.py:
from celery.utils.log import get_task_logger # 从重命名后的配置文件导入app from celery_app import app logger = get_task_logger("emailSvcCelery") class ClassTask(app.Task): name = "class-task" def run(self): print('Task Ran!') logger.info("Class task executed successfully") app.register_task(ClassTask)
3. 调整message_consumer.py的任务调用方式
Celery的Task类实例不能直接用ClassTask().delay()调用,需通过Celery实例或任务名称触发:
import os import pika # 从配置文件导入Celery app from celery_app import app connection_parameters = pika.ConnectionParameters(os.environ.get("RABBITMQ_CONN_STRING")) connection = pika.BlockingConnection(connection_parameters) channel = connection.channel() queues = ('queue-test',) for q in queues: channel.queue_declare(queue=q, durable=True) def callback(channel, method, properties, body): # 通过任务名称触发Celery任务 app.send_task("class-task") channel.basic_consume(queue='queue-test', on_message_callback=callback, auto_ack=True) print("Started Consuming...") channel.start_consuming()
验证流程
- 启动Celery Worker:
celery -A celery_app worker --loglevel=info
- 运行消息消费者脚本:
python3 message_consumer.py
- 向
queue-test队列发送消息,即可触发ClassTask执行。
内容的提问来源于stack exchange,提问作者sclem72
相关产品推荐
相关产品推荐

