Celery5.1.2调用AbortableAsyncResult的abort方法报RPC后端错误如何解决?
问题原因分析
- 后端类型不兼容:AbortableTask的实现逻辑是在结果后端存储任务的中止标记,RPC后端属于临时点对点传输组件,仅会暂存任务执行结果直到调用方获取,不支持持久化存储额外的任务元数据,所以调用
abort()时无法写入中止状态,就会抛出你遇到的RuntimeError: RPC backend missing task request异常。解决方法是将后端替换为Redis、数据库这类支持持久化的后端,比如Redis后端配置为backend = "redis://redis:6379/0"。 - 任务中止检查逻辑错误:你的代码先执行了100秒的
sleep,之后才检查is_aborted(),就算中止功能正常,任务也会先等满100秒才会响应中止请求,没有实际意义。需要把长耗时操作拆分成小片段,每执行一段就检查一次中止状态。 - AbortableTask版本说明:Celery官方确实曾经计划在4.0版本移除AbortableTask,但后续该模块作为非核心的contrib扩展被保留了下来,5.1.2版本中该组件可用,但仅支持符合要求的后端,没有官方长期维护承诺。
修正后的代码示例
from time import sleep from celery import Celery from celery.contrib.abortable import AbortableTask celery_app = Celery( broker = "amqp://login:pass@rabbitmq:5672", # 替换为支持持久化的后端 backend = "redis://redis:6379/0", include = ["root.subpath.my_tasks_module"] ) @celery_app.task(acks_late=True, bind=True, base=AbortableTask) def my_task(self, a, b): # 拆分长耗时操作,定期检查中止状态 for _ in range(100): if self.is_aborted(): return None sleep(1) return a * b def foo(): abortable_async_result = my_task.delay(4, 4) sleep(10) abortable_async_result.abort()
内容的提问来源于stack exchange,提问作者Mark Nevt
相关产品推荐
相关产品推荐

