【问题标题】:Getting all task IDs from nested chains and chords从嵌套链和和弦中获取所有任务 ID
【发布时间】:2014-05-24 07:35:57
【问题描述】:

我正在使用带有 Redis 后端的 Celery 3.1.9。我正在运行的工作由几个子任务组成,这些子任务以和弦和链的形式运行。结构如下:

  1. 准备
  2. 下载数据(2 名工作人员的和弦)
  3. 解析和存储下载的数据
  4. 4 名工人的长跑和弦
  5. 敲定
  6. 生成报告

列表中的每个项目都是一个子任务,它们都链接在一起。第 2 步和第 4 步是和弦。整个过程是通过为第 4 步创建一个和弦,其回调是 4 -> 6 的链,然后为第 2 步创建一个和弦,其回调是 3 -> 第一个和弦。然后,最后创建一个链 1 -> 第二个和弦。然后这个链以 delay() 开始,它的 ID 存储在数据库中。

问题有两个方面。首先,我希望能够撤销整个事情,其次我希望在我的 Task 类上有一个自定义的 on_failure,它会进行一些清理,并向用户报告失败。

目前我存储链的任务 ID。我想我可以用它来撤销链条。此外,如果出现错误,我想将链带到其根(在on_failure 处理程序中)以从数据库中检索相关记录。这不起作用,因为当您仅使用任务 ID 重新创建 AsyncResult 的实例时,其父属性为 None

我尝试的第二件事是将serializable() 调用的结果存储在外链的结果中。但是,这不会返回整个 AsyncResult 对象树,它只返回链中第一级的 ID(因此不是和弦中子级的 ID。)

我尝试的第三件事是自己实现serializable(),但事实证明,原始方法没有超过2级的原因是因为链的孩子是celery.canvas.chord对象,而不是@ 987654330@ 个实例。

问题说明:

chord([
    foo.si(),
    foo.si(),
    foo.si(),
], bar.si() | bar.si())
res = chord.apply_async()
pprint(res.serializable())

打印以下内容:

(('50c9eb94-7a63-49dc-b491-6fce5fed3713',
  ('d95a82b7-c107-4a2c-81eb-296dc3fb88c3',
   [(('7c72310b-afc7-4010-9de4-e64cd9d30281', None), None),
    (('2cb80041-ff29-45fe-b40c-2781b17e59dd', None), None),
    (('e85ab83d-dd44-44b5-b79a-2bbf83c4332f', None), None)])),
 None)

第一个 ID 是回调链的 ID,第二个 ID 来自 chord 任务本身,最后三个是 chord 内部的实际任务。但我无法从回调链内的任务中获取结果(即两个 bar.si() 调用的 ID)。

有什么方法可以获取实际的任务 ID?

【问题讨论】:

  • 您找到解决方案了吗?
  • 不,这是不可能的。原因是嵌套的和弦只是父任务中的一个标题,而不是队列中的一个单独的任务。当父级完成时,它会安排其他任务。
  • 完成后可以遍历和弦上的任务吗? github.com/celery/celery/blob/master/celery/canvas.py#L1242def _traverse_tasks(self, tasks, value=None):

标签: python redis celery


【解决方案1】:

一个 hacky 方法是使用 apply_async 调用任务,保存任务 ID 并手动等待它们。通过这种方式,您将完全控制发生,但您应该只等待异步任务作为最后的手段。现在您可以访问任务 id、返回值等。例如:

 task1 = a_task.apply_async()
 task2 = b_task.apply_async()
 task3 = c_task.apply_async()

 tasks = [task1, task2, task3]

 for task in tasks:
     task.wait()

【讨论】:

    【解决方案2】:

    我有一个嵌套的 dag,它是组和链的混合体。以下递归方法可以很好地获取 task_ids 及其结果:

    import celery
    
    
    
    def get_task_id_result_tuple_list(run_dag, with_result=True):
        
        task_id_result_list = []
        
        # for groups, parents are first task, then iterate over the children
        if isinstance(run_dag, celery.result.GroupResult):
            entry = (run_dag.parent, run_dag.parent.result) if with_result else run_dag.parent
            task_id_result_list.append(entry)
            children = run_dag.children
            for child in children:
                task_id_result_list.extend(get_task_id_result_tuple_list(child, with_result))
    
        # for AsyncResults, append parents in reverse
        elif isinstance(run_dag, celery.result.AsyncResult):
            ch = run_dag
            ch_list = [(ch, ch.result)] if with_result else [ch]
            while ch.parent is not None:
                ch = ch.parent
                entry = (ch, ch.result) if with_result else ch
                ch_list.append(entry)
            
            # remember to reverse the list to get the calling order
            task_id_result_list.extend(reversed(ch_list))
            
        return task_id_result_list
        
    # dag is the nested celery structure of chains and groups
    run_dag = dag.apply_async()
    
    task_id_result_tuples = get_task_id_result_tuple_list(run_dag)
    task_id_only = get_task_id_result_tuple_list(run_dag, False)
    

    注意:我还没有用和弦测试过这个,但我想它要么按原样工作,要么可能需要另一个条件分支来处理它。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2013-05-12
      • 2020-01-21
      • 2013-09-27
      • 1970-01-01
      • 1970-01-01
      • 2016-11-07
      相关资源
      最近更新 更多