【发布时间】: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