【问题标题】:Airflow: How to create child operators from inside python_callable in PythonOperator气流:如何从 PythonOperator 中的 python_callable 内部创建子运算符
【发布时间】:2020-03-07 18:50:02
【问题描述】:

我有一个简单的 python 运算符,定义如下:

loop_records = PythonOperator(
    task_id = 'loop_records',
    provide_context = True,
    python_callable = loop_topic_records,
    dag = dag    
)

这个python操作符调用loop_topic_records,定义如下:

def loop_topic_records(**context):
    parent_dag = context['dag']
    for i in range(3):
        op = DummyOperator(
            task_id="child_" + str(i),
            dag=parent_dag
        )
        logging.info('Child operator ' + str(i))
        loop_records >> op

我看到代码没有引发任何错误。它甚至在日志中打印Child operator 0..2。但是,在 dag Graph view 中我没有看到子运算符,我只看到 loop_records 节点,好像我的 dag 只包含一个运算符。那么,这有什么问题呢?我该如何解决?

【问题讨论】:

  • 我刚刚在 operatorm 上创建了它必须失败(我只是将这样的逻辑放入此运算符)。但是,当我运行整个 dag 时,它运行成功。因此,这意味着以这种方式调用的嵌套子运算符永远不会运行

标签: python airflow


【解决方案1】:

你不能做你想做的事。每个 DAG 一旦被 Airflow 加载,都是静态的,并且不能从正在运行的任务中更改。您从任务内部对 DAG 所做的任何更改都会被忽略

可以做的是启动其他 DAG,使用 airflow_multi_dagrun plugin 提供的 Multi DAG 运行运算符;创建一个 DAG 的 DAG,可以这么说:

from airflow.operators.dagrun_operator import DagRunOrder
from airflow.operators.multi_dagrun import TriggerMultiDagRunOperator

def gen_topic_records(**context):
    for i in range(3):
        # generate `DagRunOrder` objects to pass a payload (configuration)
        # to the new DAG runs.
        yield DagRunOrder(payload={"child_id": i})
        logging.info('Triggering topic_record_dag #%d', i)

loop_topic_record_dags = TriggerMultiDagRunOperator(
    task_id='loop_topic_record_dags',
    dag=dag,
    trigger_dag_id='topic_record_dag',
    python_callable=gen_topic_records,
)

以上将触发名为topic_record_dag 的 DAG 启动 3 次。在该 DAG 中的运算符内部,您可以通过dag_run.conf 对象(在模板中)或context['dag_run'].conf 引用(在PythonOperator() 代码中,设置provide_context=True)访问任何设置为有效负载的内容。

如果在这 3 个 DAG 完成后您需要做额外的工作,您只需向上述 DAG 添加一个传感器。传感器是等待直到有特定外部信息可用的操作员。在此处使用一个在所有子 DAG 完成时触发的。同一个插件有一个MultiDagRunSensor,这正是你所需要的这里,它会在TriggerMultiDagRunOperator 任务启动的所有 DAG 完成(成功或失败)时触发:

from airflow import DAG
from airflow.operators.multi_dagrun import MultiDagRunSensor

wait_for_topic_record_dags = MultiDagRunSensor(
    task_id='wait_for_topic_record_dags',
    dag=dag
)

loop_topic_record_dags >> wait_for_topic_record_dags

然后在该传感器之后放置更多运算符。

【讨论】:

    【解决方案2】:

    我不确定您在 python_callable 函数中创建子运算符的用例是什么,但如果这不是严格要求,您可以像这样循环创建 PythonOperator

    op = DummyOperator(
      task_id="dummy",
      dag=dag
    )
    
    for i in range(3):
        loop_records = PythonOperator(
            task_id=f'loop_records_{i}',
            provide_context=True,
            python_callable=loop_topic_records,
            dag=dag
        )
        loop_records >> op
    

    【讨论】:

    • 嗯,问题是,我需要在python_callable 内部进行,因为在实际情况下循环数取决于xcom。在我的问题中的loop_topic_records 中,我循环了一个静态数组,但是,实际上它是动态的并且基于xcom
    • 基本上,Airflow 做不到,任务是静态的,执行前会自动生成。但是你可以找到像这样的stackoverflow.com/questions/39133376/…
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2022-07-06
    • 2023-01-07
    • 1970-01-01
    • 2019-05-13
    • 1970-01-01
    • 2023-03-21
    • 1970-01-01
    相关资源
    最近更新 更多