【问题标题】:Chords not waiting for all child tasks和弦不等待所有子任务
【发布时间】:2013-11-17 23:48:02
【问题描述】:

我有一些这样的任务:

@celery.task
def generate():
    sleep(1.0)
    print "Generate done!"
    return 'result'

@celery.task
def lower(result):
    sleep(1.0)
    print "Lower done!"
    return result.lower()

@celery.task
def upper(result):
    sleep(1.0)
    print "Upper done!"
    return result.upper()

@celery.task
def upload(result):
    sleep(1.0)
    print "Upload done for: %s!" % (result,)
    return 'upload'

@celery.task
def callback(results):
    print "It's all done! %s" % (results,)

我正在创建一个如下所示的和弦:

chord(
    header=chain(
        generate.s(),
        group(
            chain(lower.s(), upload.s()),
            chain(upper.s(), upload.s())
        )
    ), body=callback.s()
).delay()

我遇到的问题是我的回调(应该在所有任务完成后触发)似乎在 generate 之后立即触发。

如果不清楚,工作流程是这样的:

  1. 生成一个结果,然后将其结果传递给一个组的成员,从而实现并行:
    1. 第一组将从generate获取结果,使用lower将其转换为小写,然后使用upload上传结果。
    2. 第二组将从generate 获取结果,使用upper 将其转换为大写,然后使用“上传”上传结果。
  2. 完成所有这些后,应该调用callback 任务回调。

预期

callback 任务将在启动后至少 3 秒被调用。

实际

callback 任务在启动后大约 1 秒被调用,并且不等待组成员完成执行。

以下是证明它不等待组的日志:

[2013-11-17 18:20:40,447: WARNING/PoolWorker-8] Generate done!
[2013-11-17 18:20:41,493: WARNING/PoolWorker-6] Upper done!
[2013-11-17 18:20:41,493: WARNING/PoolWorker-1] Lower done!
[2013-11-17 18:20:41,535: WARNING/PoolWorker-6] It's all done! [('e0016a35-d538-4e96-ad86-6ddf91ef4a09', [('b1af78a9-7935-4037-84e4-9fae6d7c027e', None), ('d69c4c99-af9c-476f-af7d-7f647c4d9c83', None)])]
[2013-11-17 18:20:42,522: WARNING/PoolWorker-7] Upload done for: result!
[2013-11-17 18:20:42,523: WARNING/PoolWorker-5] Upload done for: RESULT!

似乎 Celery 不等待组。有没有办法让 Celery 等到所有任务(包括组成员)执行完毕?

【问题讨论】:

    标签: celery


    【解决方案1】:

    您在这里使用链作为和弦标题,但标题必须是一个组:

    chord(
        header=chain(
            generate.s(),
            group(
                chain(lower.s(), upload.s()),
                chain(upper.s(), upload.s())
            )
        ), body=callback.s()
    ).delay()
    

    对于chain(generate.s(), group(...),没有什么可以同步的,因为 组并行发生。

    您的工作流程最好这样表达:

    filters = group(lower.s() | upload.s(),
                    upper.s() | upload.s())
    result = (generate.s() | filters | callback.s())()
    

    注意:chain(group, sig) 会自动转换为和弦

    【讨论】:

    • 我直到现在才使用链的唯一原因是我需要一个总是触发的回调,因为我的callback 任务会进行清理和报告关于什么有效,什么无效。使用上面的示例,callback 是否总是被触发?
    猜你喜欢
    • 2016-11-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-06-05
    • 2023-04-04
    • 1970-01-01
    • 2012-07-05
    • 2018-07-07
    相关资源
    最近更新 更多