【问题标题】:How to run airflow DAG with conditional tasks如何使用条件任务运行气流 DAG
【发布时间】:2020-07-08 08:04:07
【问题描述】:

总共有 6 个任务。这些任务需要根据输入 json 中一个字段的 (flag_value) 值执行。 如果 flag_value 的值为真,那么所有任务都需要以这样的方式执行, 首先task1然后并行到(task2和task3一起),并行到task4,并行到task5。 一旦这一切完成,然后任务6。 由于是气流和 DAG 的新手,我不知道如何在这种情况下运行。

如果 flag_value 的值为 false,则顺序仅是顺序的
任务_1 >> 任务_4 >> 任务_5 >> 任务_6。

下面是我的 DAG 代码。

from airflow import DAG
from datetime import datetime
from airflow.providers.databricks.operators.databricks import DatabricksSubmitRunOperator


default_args = {
    'owner': 'airflow',
    'depends_on_past': False
}

dag = DAG('DAG_FOR_TEST',default_args=default_args,schedule_interval=None,max_active_runs=3, start_date=datetime(2020, 7, 8)) 


#################### CREATE TASK #####################################   

task_1 = DatabricksSubmitRunOperator(
    task_id='task_1',
    databricks_conn_id='connection_id_details',
    existing_cluster_id='{{ dag_run.conf.clusterId }}',
    libraries= [
        {
        'jar': 'dbfs:/task_1/task_1.jar'
        }        
        ],
    spark_jar_task={
        'main_class_name': 'com.task_1.driver.TestClass1',
        'parameters' : [
            '{{ dag_run.conf.json }}'       
        ]
    }
)



    
task_2 = DatabricksSubmitRunOperator(
    task_id='task_2',
    databricks_conn_id='connection_id_details',
    existing_cluster_id='{{ dag_run.conf.clusterId }}',   
    libraries= [
        {
        'jar': 'dbfs:/task_2/task_2.jar'
        }        
        ],
    spark_jar_task={
        'main_class_name': 'com.task_2.driver.TestClass2',
        'parameters' : [
            '{{ dag_run.conf.json }}'                               
        ]
    }
)
    
task_3 = DatabricksSubmitRunOperator(
    task_id='task_3',
    databricks_conn_id='connection_id_details',
    existing_cluster_id='{{ dag_run.conf.clusterId }}',   
    libraries= [
        {
        'jar': 'dbfs:/task_3/task_3.jar'
        }        
        ],
    spark_jar_task={
        'main_class_name': 'com.task_3.driver.TestClass3',
        'parameters' : [
            '{{ dag_run.conf.json }}'   
        ]
    }
) 

task_4 = DatabricksSubmitRunOperator(
    task_id='task_4',
    databricks_conn_id='connection_id_details',
    existing_cluster_id='{{ dag_run.conf.clusterId }}',
    libraries= [
        {
        'jar': 'dbfs:/task_4/task_4.jar'
        }        
        ],
    spark_jar_task={
        'main_class_name': 'com.task_4.driver.TestClass4',
        'parameters' : [
            '{{ dag_run.conf.json }}'   
        ]
    }
) 

task_5 = DatabricksSubmitRunOperator(
    task_id='task_5',
    databricks_conn_id='connection_id_details',
    existing_cluster_id='{{ dag_run.conf.clusterId }}',
    libraries= [
        {
        'jar': 'dbfs:/task_5/task_5.jar'
        }        
        ],
    spark_jar_task={
        'main_class_name': 'com.task_5.driver.TestClass5',
        'parameters' : [
            'json ={{ dag_run.conf.json }}' 
        ]
    }
) 

task_6 = DatabricksSubmitRunOperator(
    task_id='task_6',
    databricks_conn_id='connection_id_details',
    existing_cluster_id='{{ dag_run.conf.clusterId }}',
    libraries= [
        {
        'jar': 'dbfs:/task_6/task_6.jar'
        }        
        ],
    spark_jar_task={
        'main_class_name': 'com.task_6.driver.TestClass6',
        'parameters' : ['{{ dag_run.conf.json }}'   
        ]
    }
) 


flag_value='{{ dag_run.conf.json.flag_value }}'

#################### ORDER OF OPERATORS ###########################  

if flag_value == 'true':
    
    task_1.dag = dag
    task_2.dag = dag
    task_3.dag = dag
    task_4.dag = dag
    task_5.dag = dag
    task_6.dag = dag
    
    task_1  >> [task_2 , task_3] >> [task_4] >> [task_5]  >> task_6    // Not sure correct 
else:
    task_1.dag = dag
    task_4.dag = dag
    task_5.dag = dag
    task_6.dag = dag
    
    task_1 >> task_4 >> task_5 >> task_6

    

        
        
    

【问题讨论】:

    标签: python airflow directed-acyclic-graphs


    【解决方案1】:

    首先,依赖不正确,这应该可以工作:

    task_1 >> [task_2 , task_3] >> task_4 >> task_5  >> task_6
    

    无法使用list_1 >> list_2 对任务进行排序,但有一些辅助方法可以提供此功能,请参阅:cross_downstream

    对于分支,您可以使用BranchPythonOperator 更改您的任务的trigger rules。不确定下面的代码,它可能有小错误,但这里的想法有效。

    task_4.trigger_rule = "none_failed"
    
    dummy = DummyOperator(task_id="dummy", dag=dag)
    
    branch = BranchPythonOperator(
        task_id="branch",
        # jinja template returns string "True" or "False"
        python_callable=lambda f: ["task_2" , "task_3"] if f == "True" else "dummy",
        op_kwargs={"f": flag_value},
        dag=dag)
    
    task_1 >> branch
    branch >> [task_2 , task_3, dummy] >> task_4
    task_4 >> task_5 >> task_6
    

    可能有更好的方法来做到这一点。

    【讨论】:

    • 非常感谢@mustafagok 的回答。但我需要的是第一个 task1 需要运行。之后 4 个任务需要以并行方式运行(task2 和 task3 一起),并行到 task4,并行到 task5。一旦这一切完成,task6 就需要运行了。
    • 第一个 task_1 正在运行。之后就不明白了但是您可以轻松更改顺序或操作,只需尝试并查看图形视图,如果不正确,请更改。如果您通过绘图或更详细地描述,我可以尝试提供帮助。 2-3-4和5都是平行的吗?
    • Task1 并行 ----> [Task2 &Task3] 并行 ----> Task4 并行 ----> Task5 全部执行完Task6
    • 不确定,与箭头平行是什么意思。 if 2-3-4-5 parallel: task_1 >> branch >> [task_2, task_3, task_4, task_5] >> task_6 if (2&3)-4-5 parallel: task_1 >> branch >> [task_2, task_4, task_5] task_2 >> task_3 [task_3, task_4, task_5] >> task_6 如果你尝试不同的版本,你会找到你需要的图表。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-05-03
    • 1970-01-01
    • 2022-11-02
    • 2019-07-21
    • 1970-01-01
    • 2017-01-01
    • 1970-01-01
    相关资源
    最近更新 更多