【发布时间】:2021-07-26 13:05:04
【问题描述】:
我们正在设计一个变量选择和参数设置器逻辑,当 DAG 被触发时需要对其进行评估。我们的 DAG 是在执行之前生成的。我们决定将我们的静态代码修改为自定义宏。
直到此时,在运算符定义之间定义了一个代码,因此当 DAG 由 DAG 生成器代码生成时,它正在运行。此代码无法处理用于选择正确 Airflow 变量的运行时参数。
for table_name in ast.literal_eval(Variable.get('PYTHON_LIST_OF_TABLES')):
dag_id = "TableLoader_" + str(table_name)
default_dag_args={...}
schedule = None
globals()[dag_id] = create_dag(dag_id, schedule, default_dag_args)
def create_dag(dag_id, schedule, default_dag_args):
with DAG(
default_args=default_dag_args,
dag_id=dag_id,
schedule_interval=schedule,
user_defined_macros={ "load_type_handler": load_type_handler }
) as dag:
# static python code which sets pipeline_rt_args for all generated DAGs the same way
# this static code could set only one type (INITIAL or INCREMENTAL)
# but we want to decide during the execution now
# Operator Definitions
OP_K = CloudDataFusionStartPipelineOperator(
task_id='PREFIX_'+str(table),
# ---> Can't handle runtime parameters <---
runtime_args=pipeline_rt_args,
# ...
)
OP_1 >> OP_2 >> ... >> OP_K >> ... >> OP_N
return dag
现在我们想要在从 UI 或 REST API 触发 DAG 时传递 load_type(例如:INITIAL,INCREMENTAL),因此我们需要修改这个旧的(静态)行为(它只处理一种情况,但不是两种情况)来获取正确的气流变量并为我们的CloudDataFusionStartPipelineOperator创建正确的对象:
例如:
{"load_type":"INCREMENTAL"}
# or
{"load_type":"INITIAL"}
但如果我们这样做:
def create_dag(dag_id, schedule, default_dag_args):
def extend_runtime_args(prefix, param, field_name, date, job_id):
# reading the Trigger-time parameter
load_type = param.conf["load_type"]
# getting the proper Airflow Variable (depending on current load type)
result = eval(Variable.get(prefix+'_'+load_type+'_'+dag_id))[field_name]
# setting 'job_id', 'dateValue', 'date', 'GCS_Input_Path' for CloudDataFusionStartPipelineOperator
# ...
return rt_args
with DAG( #...
user_defined_macros={
"extend_runtime_args": extend_runtime_args
}) as dag:
# removed static code (which executes only in generation time)
# Operator Definitions
OP_K = CloudDataFusionStartPipelineOperator(
task_id='PREFIX_'+str(table),
# ---> handles runtime arguments with custom macro <---
runtime_args="""{{ extend_runtime_args('PREFIX', dag_run, 'runtime_args', macros.ds_format(yesterday_ds_nodash,"%Y%m%d","%Y_%m_%d"), ti.job_id) }}""",
# ...
)
OP_1 >> OP_2 >> ... >> OP_K >> ... >> OP_N
return dag
注意:
这里我们需要的是自定义逻辑的“未来”评估(不在 DAG 生成时评估),它将返回一个对象,这就是我们在这里使用模板的原因。
我们会遇到以下情况:
- 在自定义宏函数
extend_runtime_args内部,返回类型是一个对象 - 评估 Jinja 模板后,返回类型更改为字符串
-
CloudDataFusionStartPipelineOperator失败,因为runtime_args属性是字符串而不是对象
问题:
- 我们如何在评估 Jinja 模板后返回一个对象(并在“未来”这样做)?
- 我们能以某种方式转换字符串吗?
- 如何保证这里的逻辑是在 DAG 执行后执行,而不是在 DAG 生成后执行?
- Jinja 模板/自定义宏在这里处理触发时间参数是好还是坏?
【问题讨论】:
-
从 Airflow 2.1 开始,您可以使用参数
render_template_as_native_obj将模板的值作为 python 对象返回。见stackoverflow.com/questions/64235496/…
标签: python jinja2 airflow google-cloud-data-fusion