【问题标题】:Celery chord executed before all tasks are completed在所有任务完成之前执行的 Celery chord
【发布时间】:2014-09-26 14:42:31
【问题描述】:

我有一堆可以同时执行的任务,但是一旦一切准备就绪,我想执行最后一个。我正在使用以下代码:

chunk_tasks = []
for index, chunk in enumerate(chunks):
    chunk_tasks.append(import_chunk.s(meta.pk))

g = group(chunk_tasks)
chord(g)(import_completed.s(meta.pk, max_lines=max_lines))

但是看起来import_completed 在所有任务完成之前执行。 import_chunk 任务也看起来像:

@task(bind=True, ignore_result=IGNORE_RESULTS)
def import_chunk(self, meta_pk):
    try:
        # do some stuff
    except Exception, e:
        if self.max_retries == self.request.retries:
            logger.exception('Unexpected error in import_chunk')
        raise self.retry(countdown=60, max_retries=3)

所以问题是我做错了什么?

【问题讨论】:

    标签: python django celery django-celery


    【解决方案1】:

    Chord 是一项仅在组中的所有任务都执行完后才执行的任务。因此,它需要任务状态在其标头中进行同步。

    但是当您将ignore_result 设置为您的task 时,worker 不会存储任务状态并为此任务返回值。

    根据您的工作流程,这将导致重试任务或抛出异常或任何故障。

    所以,chord(add.s(i, i) for i in range(10))(tsum.s()).get() 是完全有效的,并为案例 1 提供了结果,但给案例 2 带来了一些麻烦。

    案例 1:

    @app.task
    def add(x, y):
        return x + y
    
    @app.task
    def tsum(numbers):
        return sum(numbers)
    

    案例 2:

    @app.task(ignore_result=True)
    def add(x, y):
        return x + y
    
    @app.task(ignore_result=True)
    def tsum(numbers):
        return sum(numbers)
    

    因此,您必须更改 ignore_result 或更改任务的工作流程。

    来自文档:

    你应该尽可能避免使用和弦。尽管如此,和弦仍然是您工具箱中的强大原语,因为同步是许多并行算法的必需步骤。

    【讨论】:

    • 是的,我已经偶然发现了 ignore_result 的问题,它被设置为 False。
    猜你喜欢
    • 2013-03-06
    • 1970-01-01
    • 2018-01-02
    • 1970-01-01
    • 2013-09-17
    • 1970-01-01
    • 1970-01-01
    • 2017-07-18
    • 2017-06-27
    相关资源
    最近更新 更多