【问题标题】:Celery apply_async pass kwargs to all tasks in chainCelery apply_async 将 kwargs 传递给链中的所有任务
【发布时间】:2019-11-26 04:09:36
【问题描述】:

一个 Celery 任务队列,用于计算 (2 + 2) - 3 的结果。

@app.task()
def add(**kwargs):
    time.sleep(5)
    x, y = kwargs['add'][0], kwargs['add'][1]
    return x + y

@app.task()
def sub(**kwargs):
    time.sleep(5)
    x = args[0]
    y = kwargs['sub'][0]
    return x - y

示例任务数据 = kwargs = {'add' : (2, 2), 'sub' : (3)}
链接任务:result = (add.s() | sub.s()).apply_async(kwargs = kwargs)

As per design, apply_async 仅将 kwargs 应用于链中的第一个任务。我需要改变什么才能达到预期的结果?

【问题讨论】:

  • 链接任务将应用其父任务的结果作为第一个参数,因此您可以修改return以包含kwargs
  • 实际上,任务将执行非常不同的事情。例如 - task1 - 打开浏览器到url。 task2 - 将data 发送到特定的input 字段。你可以看到从task1返回东西不是一个好的设计。
  • 有人可以帮忙吗?

标签: python celery celery-task


【解决方案1】:

因此,从 Celery v4.4.0rc4 开始,除了将 kwargs 传递给每个任务的签名之外,没有更好的方法可以做到这一点。虽然它看起来像Ask Solem (Celery dev) is open to a feature request.。

链应该是这样的:

result = (add.s(job_data = job_data)| sub.s(job_data = job_data)).apply_async()

但是,由于我们的链有 10 多个任务,我不得不想出一个更简单的方法来编写这个。

# Workflow generator
def workflow_generator(task_list, job_data):
    _tasks = tuple(getattr(task, 's')(job_data = job_data) for task in task_list)
    return chain(*_tasks).apply_async()

taskList = [add, sub]
job_data = {'add' : (2, 2), 'sub' : (3)}
result = workflow_generator(taskList, job_data) 

【讨论】:

    【解决方案2】:

    为什么不只是在链接之前部分绑定任务?

    result = (add.s(add=(2,2)) | sub.s(sub=(3,3))).apply_async()
    

    【讨论】:

      猜你喜欢
      • 2012-03-22
      • 2016-02-28
      • 2013-11-23
      • 2018-10-01
      • 2021-11-22
      • 2014-03-07
      • 2018-10-19
      • 2018-01-04
      • 2015-04-12
      相关资源
      最近更新 更多