【发布时间】:2017-08-13 07:15:35
【问题描述】:
我正在努力熟悉 Airflow 并喜欢它。
然而,我有点不清楚的是如何正确参数化我的 dag,我想在哪里运行相同的 dag,但同时为多个业务线 (lob) 运行。所以基本上我想在每次运行中为多个 lob 运行下面的 dag,并让每个 lob 并行运行。
假设我定义了一个变量,它是一个 lob 数组,如“lob1”、“lob2”等。我想将下面 bigquery sql 语句中的“mylob”替换为“lob1”,然后是“lob2”等.
我在想也许我可以将 lobs 存储为 ui 中的变量,然后在 dag 中循环,但我不确定这是否最终会是连续的,因为它等待每个任务在每个任务中完成循环迭代。
我认为另一种方法可能是将这个参数化的 dag 用作更大驱动程序 dag 中的子 dag。但再次不确定这是否是最佳实践方法。
非常感谢任何帮助或指示。我觉得我在这里遗漏了一些明显的东西,但在任何地方都找不到这样的例子。
"""
### My first dag to play around with bigquery and gcp stuff.
"""
from airflow import DAG
from airflow.operators.bash_operator import BashOperator
from datetime import datetime, timedelta
from dateutil import tz
from airflow.contrib.hooks.bigquery_hook import BigQueryHook
from airflow.contrib.operators.bigquery_operator import BigQueryOperator
default_args = {
'owner': 'airflow',
'depends_on_past': False,
'start_date': datetime(2017, 3, 10),
'email': ['xxx@xxx.com'],
'email_on_failure': True,
'retries': 1,
'retry_delay': timedelta(minutes=5),
# 'queue': 'bash_queue',
# 'pool': 'backfill',
# 'priority_weight': 10,
# 'end_date': datetime(2016, 1, 1),
}
with DAG('my_bq_dag_2', schedule_interval='30 */1 * * *',
default_args=default_args) as dag:
bq_msg_1 = BigQueryOperator(
task_id='my_bq_task_1',
bql='select "mylob" as lob, "Hello World!" as msg',
destination_dataset_table='airflow.test1',
write_disposition='WRITE_TRUNCATE',
bigquery_conn_id='gcp_smoke'
)
bq_msg_1.doc_md = """\
#### Task Documentation
Append a "Hello World!" message string to the table [airflow.msg]
"""
bq_msg_2 = BigQueryOperator(
task_id='my_bq_task_2',
bql='select "mylob" as lob, "Goodbye World!" as msg',
destination_dataset_table='airflow.test1',
write_disposition='WRITE_APPEND',
bigquery_conn_id='gcp_smoke'
)
bq_msg_2.doc_md = """\
#### Task Documentation
Append a "Goodbye World!" message string to the table [airflow.msg]
"""
# set dependencies
bq_msg_2.set_upstream(bq_msg_1)
更新:试图让这个工作,但它似乎永远无法进入 lob2
"""
### My first dag to play around with bigquery and gcp stuff.
"""
from airflow import DAG
from airflow.operators.bash_operator import BashOperator
from datetime import datetime, timedelta
from airflow.contrib.operators.bigquery_operator import BigQueryOperator
default_args = {
'owner': 'airflow',
'depends_on_past': False,
'start_date': datetime(2017, 3, 10),
'email': ['xxx@xxx.com'],
'email_on_failure': True,
'retries': 1,
'retry_delay': timedelta(minutes=5),
}
dag = DAG('my_bq_dag_2', schedule_interval='@once',default_args=default_args)
lobs = ["lob1","lob2","lob3"]
for lob in lobs:
templated_command = """
select '{{ params.lob }}' as lob, concat(string(current_timestamp()),' - Hello - {{ ds }}') as msg
"""
bq_msg_1 = BigQueryOperator(
dag = dag,
task_id='my_bq_task_1',
bql=templated_command,
params={'lob': lob},
destination_dataset_table='airflow.test1',
write_disposition='WRITE_APPEND',
bigquery_conn_id='gcp_smoke'
)
【问题讨论】:
-
似乎这个例子可能有点像我需要使用分支github.com/apache/incubator-airflow/blob/master/airflow/contrib/…
-
这里的另一个选项使用 conf 参数 github.com/andrewm4894/random/blob/master/my_bq_dag_2.py