【问题标题】:Celery: How to Wait nested tasks finish in group芹菜:如何等待嵌套任务在组中完成
【发布时间】:2020-10-11 04:34:57
【问题描述】:

我需要使用一些参数向 API 发出请求,并一一处理来自请求的数据。但是有时响应数据是分页的,这意味着我需要使用相同的参数发出额外的请求。 Celerygroup让我们可以一个一个的运行任务,但是如果任务产生子任务,子任务又可以产生更多的子任务……,

有没有办法在运行组中的下一个任务之前等待所有子任务完成?或者 celery 给了我们更好的方法来解决我的任务?

def some_api_call(start_date=None, end_date=None, token_value=None):
    pass


@celery_app.task(name='run_task', max_retries=None, bind=True)
def run_task(self):
    group_items = [
        task1.s('2020-01-02', '2020-01-03'),
        task1.s('2020-01-05', '2020-01-01'),
        task1.s('2020-01-010', '2020-01-04'),
    ]
    group(group_items)()

@celery_app.task(name='task1', max_retries=None, bind=True)
def task1(self, start_date=None, end_date=None, token_value=None, *args, **kwargs):
    res = some_api_call(start_date, end_date, token_value)
    if res['token_value']:
        # NEXT ELEMENT IN THE GROUP SHOULD WAIT UNTIL NESTED CHILD TASKS DONE
        task1.delay(token_value=token_value)

经纪人 - Redis。 我可能的解决方案[伪代码]:

  1. 等到父任务中的子任务完成
    res = task1.delay(token_value=token_value) 

    res.get()

解决方案不好 - 我们挡住了胎面。 不确定 celery 是否有 await 替代品。

  1. 用户task.retry()检查子任务是否完成。
    taskid = task1.delay(token_value=token_value)
    if AsyncResult(taskid).state != "successes": 

      self.retry()

所以我们将重试父任务,直到子任务完成并且不要阻塞线程。 但就像在第一个解决方案中一样:父任务处理其数据,但状态会重试。

【问题讨论】:

    标签: python celery django-celery


    【解决方案1】:

    如果您需要等待组中的任务完成,您应该使用Chord 原语。此外,如果您需要按顺序执行任务(一个接一个),请使用 Chain 原语。 Chord 基本上是一个 Group 链,也是一个最终的任务……

    【讨论】:

    • 你能告诉我如何在我的情况下使用和弦吗?如果您阅读了我的问题 - 我需要等待所有子任务完成,然后再运行组中的下一个任务。
    • Chord 基本上是一个 Group + 任务,在 Group 完成后“收集”数据。这一切都在文档中得到了很好的解释。如果您想按顺序运行任务,还有另一个用于此目的的原语 - Chain。
    • Chord 不能解决我的问题,chord 只是组完成时的回调任务。我需要知道组中的任务是否完成了所有子任务
    • 然后你制作一个稍微复杂一点的工作流程,其中你的“子任务”也是和弦......
    猜你喜欢
    • 2020-05-19
    • 2020-07-09
    • 2019-06-12
    • 1970-01-01
    • 1970-01-01
    • 2014-10-28
    • 2015-08-15
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多