【发布时间】:2016-08-30 14:53:51
【问题描述】:
我有一个 Celery 任务,它从 SQS 队列中获取消息并尝试运行它。如果失败,它应该每 10 秒重试至少 144 次。我认为正在发生的事情是它失败并重新进入队列,同时它创建一个新的,将其复制到 2。这 2 再次失败并遵循相同的模式创建 2 个新的并成为 4 个消息全部的。所以如果我让它运行一段时间,队列就会被阻塞。
我没有得到的是在不重复的情况下重试的正确方法。以下是重试的代码。请看看是否有人可以在这里指导我。
from celery import shared_task
from celery.exceptions import MaxRetriesExceededError
@shared_task
def send_br_update(bgc_id, xref_id, user_id, event):
from myapp.models.mappings import BGC
try:
bgc = BGC.objects.get(pk=bgc_id)
return bgc.send_br_update(user_id, event)
except BGC.DoesNotExist:
pass
except MaxRetriesExceededError:
pass
except Exception as exc:
# retry every 10 minutes for at least 24 hours
raise send_br_update.retry(exc=exc, countdown=600, max_retries=144)
更新: 问题的更多解释...
用户在我的数据库中创建了一个对象。其他用户对该对象采取行动,当他们更改该对象的状态时,我的代码会发出信号。然后信号处理程序启动一个 celery 任务,这意味着它连接到所需的 SQS 队列并将消息提交到队列。运行工作程序的 celery 服务器看到该新消息并尝试执行任务。这是它失败的地方,重试逻辑进来了。
根据retry a task 的 celery 文档,我们需要做的就是使用 countdown 和/或 max_retries 引发 self.retry() 调用。如果 celery 任务引发异常,则将其视为失败。我不确定 SQS 如何处理这个问题。我所知道的是一个任务失败,队列中有两个,这两个都失败了,然后队列中有 4 个,依此类推......
【问题讨论】:
-
任务看起来没问题,请添加运行任务的代码。
-
这是标准的 celery 任务。由 celery worker 在单独的服务器上运行。在我的例子中,代理是 Amazon SQS。
标签: python django task celery amazon-sqs