【问题标题】:Celery chaining executing subtask before group tasks completed芹菜链在组任务完成之前执行子任务
【发布时间】:2014-10-28 15:09:58
【问题描述】:

我要做的是执行一组任务,并将组中每个任务的结果发送到另一个任务,然后将所有这些任务的结果发送到一个最终任务。

例如

jobgroup = (
           group(tasks.task1.s(p) for p in params) |
           tasks.dmap.s(tasks.task2.s())
          )

这部分对我来说没问题,但我遇到的问题是,如果我试图将所有结果放入一个最终任务中,例如

mychain = chain(jobgroup | tasks.task3.s())

我看到的是 task3() 正在被调用,而来自 task2 的任务仍处于挂起状态,方法是在它们进入时打印出状态(以下 sn-p 是我在 task3 函数中打印的内容:

@task()
def task3(input):
   for item in input.results:
      logger.info(item.status)

记录结果

[2014-09-04 10:48:41,905: INFO/Worker-7] tasks.task3[27053688-3c5c-4ca5-975f-356a66d55364]: PENDING
[2014-09-04 10:48:41,905: INFO/Worker-7] tasks.task3[27053688-3c5c-4ca5-975f-356a66d55364]: PENDING

那么我该如何设置,以便在作业组中的所有任务完成之前不调用 task3?

【问题讨论】:

  • 我在这里也注意到的其他一点是父任务完成并按预期返回所有结果但是我的任务3仍然有待处理的项目

标签: python celery


【解决方案1】:

Chord:

A chord consists of a header and a body. The header is a group of tasks that MUST COMPLETE before the callback is called. A chord is essentially a callback for a group of tasks.

示例:

>>> from celery import chord
>>> res = chord([add.s(2, 2), add.s(4, 4)])(sum_task.s())

您可以使用chord 来完成任务。

from celery import chord, group, chain

task_1 = group(tasks.task1.s(p) for p in params)
task_2 = group(tasks.dmap.s(tasks.task2.s())
task_3 = tasks.task3.s()

workflow = chord(chain([task_1, task_2]))(task_3)

【讨论】:

  • Chord 抛出这个: Traceback(最近一次调用最后一次):文件“run.py”,第 108 行,在 工作流 = chord([task_1, task_2])(task_3) File " /Library/Python/2.7/site-packages/celery/canvas.py",第 611 行,在 call 中 return self.apply_async((), {'body': body} if body else {} , **options) 文件“/Library/Python/2.7/site-packages/celery/canvas.py”,第 602 行,在 apply_async _chord = self.type 文件“/Library/Python/2.7/site-packages/celery/ canvas.py",第 590 行,类型为 app = self.tasks[0].type.app AttributeError: 'generator' object has no attribute 'type'
  • 只是为了确保这里是完整的示例应用程序:gist.github.com/sinkers/d4e0bcac5f226e8bcc21 >celery --version 3.1.13 (Cipater)
  • @Andrew 因为 task_1 是任务列表,所以必须使用 groupjob。修改了代码。可以再试一次吗?
  • 我现在在工作日志中收到错误:[2014-09-06 16:32:23,142: ERROR/Worker-2] 弦回调“1059e974-ad4c-4177-a894-fc410b9edd7b”引发:ValueError('1059e974-ad4c-4177-a894-fc410b9edd7b',) Traceback(最近一次调用最后):文件“/Library/Python/2.7/site-packages/celery/backends/base.py”,第 543 行,在on_chord_part_return raise ValueError(gid) ValueError: 1059e974-ad4c-4177-a894-fc410b9edd7b
  • @Andrew task_2 & task_1 必须链接在一起。它现在可以工作了。
猜你喜欢
  • 2020-05-19
  • 2020-08-07
  • 2014-12-04
  • 1970-01-01
  • 2020-07-09
  • 2021-12-16
  • 2018-03-07
  • 1970-01-01
  • 2015-04-29
相关资源
最近更新 更多