【问题标题】:How can we use Deferred DAG assignment in airflow我们如何在气流中使用延迟 DAG 分配
【发布时间】:2020-07-01 09:47:38
【问题描述】:

我是 Apache 气流新手并使用 DAG。 y代码如下。

在输入 json 中,我有一个名为“sports_category”的参数。如果其值为'football',则需要运行 football_players 任务,如果其值为 cricket,则需要运行'cricket_players' 任务。

import airflow

from airflow import DAG
from airflow.contrib.operators.databricks_operator import DatabricksSubmitRunOperator

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2020, 6, 23)
}

dag = DAG('PLAYERS_DETAILS',default_args=default_args,schedule_interval=None,max_active_runs=5) 

football_players = DatabricksSubmitRunOperator(
    task_id='football_players',
    databricks_conn_id='football_players_details',
    existing_cluster_id='{{ dag_run.conf.clusterId }}',
    libraries= [
        {
        'jar': {{ jar path }}
        }        
        ],
        databricks_retry_limit = 3,
    spark_jar_task={
        'main_class_name': 'football class name1',
        'parameters' : [
            'json ={{ dag_run.conf.json }}'     
        ]
    }
)

cricket_players = DatabricksSubmitRunOperator(
    task_id='cricket_players',
    databricks_conn_id='cricket_players_details',
    existing_cluster_id='{{ dag_run.conf.clusterId }}',
    libraries= [
        {
        'jar': {{ jar path }}
        }        
        ],
        databricks_retry_limit = 3,
    spark_jar_task={
        'main_class_name': 'cricket class name2',
        'parameters' : [
            'json ={{ dag_run.conf.json }}'     
        ]
    }
)

【问题讨论】:

    标签: airflow directed-acyclic-graphs


    【解决方案1】:

    我建议使用 BranchPythonOperator,它将一个函数作为参数,并根据函数内部编写的逻辑返回流程中下一个需要运行的 TASK。

    请参阅 here 获取文档和 here 获取示例 dag。

    让我知道你的回应!

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-03-31
      • 2022-11-10
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多