如何在Airflow中从Jinja模板字符串获取Python字典或列表作为算子参数
核心问题原因
Airflow 所有 Jinja 模板的渲染逻辑默认在任务执行阶段才执行,DAG 解析(算子实例化)阶段你写的{{ xxx }}只是一个未被渲染的字符串。另外默认渲染结果会被转为字符串类型,哪怕你拉取的XCom或者宏返回的是原生列表/字典,默认也会被序列化为字符串,自然无法直接作为需要原生类型的算子属性传入。
Airflow 1.x 解决方案
方案1:修改自定义算子的模板字段(最优解,适配自定义Operator场景)
所有Airflow Operator都有template_fields类属性,列表中声明的属性会在任务执行前自动完成模板渲染。你只需要两步:
- 在你的
ExampleOperator中把需要动态传值的属性加入template_fields - 在算子的
execute方法中对渲染后的结果做类型转换,推荐用ast.literal_eval(比原生eval更安全,不会执行恶意代码)
代码示例:
from airflow.models import BaseOperator import ast class ExampleOperator(BaseOperator): # 把需要支持模板渲染的属性加到这里 template_fields = ("property_needs_list", "property_needs_dict",) def __init__(self, property_needs_list, property_needs_dict, *args, **kwargs): super().__init__(*args, **kwargs) self.property_needs_list = property_needs_list self.property_needs_dict = property_needs_dict def execute(self, context): # 执行到这里的时候,上面两个属性已经完成模板渲染 if isinstance(self.property_needs_list, str): self.property_needs_list = ast.literal_eval(self.property_needs_list) if isinstance(self.property_needs_dict, str): self.property_needs_dict = ast.literal_eval(self.property_needs_dict) # 后续业务逻辑直接用这两个原生类型属性即可 print(type(self.property_needs_list)) # 输出 <class 'list'> print(type(self.property_needs_dict)) # 输出 <class 'dict'>
之后你原来的模板写法直接就能生效:
doExampleTask = ExampleOperator( task_id = "doExampleTask", property_needs_list = '{{ ti.xcom_pull(task_ids="PreviousTask", key="list_structure") }}', property_needs_dict = '{{ generate_dict() }}', )
方案2:用PythonOperator封装(适配无法修改Operator源码的场景)
如果ExampleOperator是第三方提供、你没法修改源码,就用PythonOperator在执行阶段动态调用它:
from airflow.operators.python_operator import PythonOperator def run_example_op(**context): # 执行阶段直接拉取XCom/调用宏,拿到的就是原生类型 list_val = context["ti"].xcom_pull(task_ids="PreviousTask", key="list_structure") dict_val = generate_dict() # 直接调用你的自定义宏函数即可 # 动态实例化算子并执行 op = ExampleOperator( task_id="inner_do_example", property_needs_list=list_val, property_needs_dict=dict_val ) return op.execute(context) doExampleTask = PythonOperator( task_id="doExampleTask", python_callable=run_example_op, provide_context=True )
Airflow 2.x 解决方案
2.x版本原生提供了原生类型渲染支持,不需要改Operator代码,只需要在DAG实例化的时候开启render_template_as_native_obj=True即可,渲染后的模板结果会直接保留Python原生类型:
from airflow import DAG from datetime import datetime # 开启原生类型渲染 with DAG( dag_id="your_dag_id", start_date=datetime(2023,1,1), render_template_as_native_obj=True, # 核心参数 schedule_interval=None ) as dag: # 你原来的写法直接就能用,渲染后自动就是列表/字典类型 doExampleTask = ExampleOperator( task_id = "doExampleTask", property_needs_list = '{{ ti.xcom_pull(task_ids="PreviousTask", key="list_structure") }}', property_needs_dict = '{{ generate_dict() }}', )
如果需要兼容旧版逻辑,也可以继续使用Airflow 1.x的两种方案,2.x版本完全兼容。
内容的提问来源于stack exchange,提问作者elaspog
相关产品推荐
相关产品推荐

