【问题标题】:Task progress is not updated latest status on Celery+RabbitMQ任务进度未更新 Celery+RabbitMQ 最新状态
【发布时间】:2015-03-05 05:29:10
【问题描述】:

我在 Celery + RabbitMQ 结果后端使用custom states 实现了长任务的进度反馈。

但是调用者无法像我预期的那样检索最新的进度状态。在下面的代码中,result.info['step']总是返回0,然后任务将以“result=42”结束。

# tasks.py -- celery worker
from celery import Celery
app = Celery('tasks', backend='amqp', broker='amqp://guest@localhost//')

@app.task
def long_task():
  for i in range(0, 10):
    timer.sleep(10)  # some work
    self.update_state(state='PROGRESS', meta={'step': i})
  return 42


# caller.py
from tasks import long_task
result = long_task.delay()

while not (result.successful() or result.failed()):
  try:
    result.get(timeout=1)
  except celery.exceptions.TimeoutError:
    if result.state == 'PROGRESS':
      print("progress={}".format(result.info['step']))
print("result={}".format(result.get()))

Python 3.4.1 / Celery 3.1.17 / RabbitMQ 3.4.4

【问题讨论】:

    标签: python rabbitmq celery


    【解决方案1】:

    我认为这是一个微妙的时间问题,再加上RabbitMQ result backend 将任务结果作为消息发送并且只能检索一次这一事实。

    预先简短回答:在您真正需要最终结果之前避免致电result.get():

    while not result.ready():
      if result.state == "PROGRESS":
        print("progress={}".format(result.info['step']))
      time.sleep(1)
    print("result={}".format(result.get()))
    # +additional cleanup: see comments below
    

    更长的答案是,这里实际上有两种方法(和一个属性)在与 AMQP 后端对话:

    • AsyncResult.get()

      调用AMQPBackend.wait_for(),它会消耗任务队列中的所有结果,直到出现状态为celery.states.READY_STATES 的结果。

    • AsyncResult.successful()、AsyncResult.failed()、AsyncResult.info

      调用AMQPBackend.get_task_meta(),它消耗队列中任务的所有结果,然后缓存并返回最新的结果。如果未检索到任何消息,则后端返回缓存结果或PENDING 结果。注意:后端最新消息是requeued,如果是the final result,会被AsyncResult实例缓存1。

    调用result.get() 将消耗所有状态更新,result.info 没有机会提供最新的进度报告;相反,它很可能是一个陈旧的缓存,其中一个对AsyncResult.get_task_meta() 的调用在某个时候设法获取了它。

    因此,根据时间的不同,step 可能会在次最坏的情况下卡在 0,其中最坏的情况是 PROGRESS 状态永远不会到达调用者。

    1由于最终结果在通过调用get_task_meta() 获取时会被重新排队和缓存,因此您需要手动排空队列,如下面的评论中所述。

    【讨论】:

    • 问题解决了!但它会导致另一个问题,即 RabbitMQ 不会因为未使用的最终“SUCCESS”消息而自动删除“结果队列”。这是肮脏的临时解决方法:result.backend.wait_for(task_id=result.task_id, cache=False).
    • 我已经更新了关于为什么最后一条消息留在队列中的答案。这可能是意外行为。
    猜你喜欢
    • 2019-08-15
    • 1970-01-01
    • 2018-08-21
    • 1970-01-01
    • 2018-05-24
    • 2012-06-21
    • 1970-01-01
    • 2015-09-12
    • 1970-01-01
    相关资源
    最近更新 更多