如何通过Airflow自定义Operator将数据传入SQL模板?
解决Airflow自定义Operator传递参数渲染SQL的问题
你的问题根源在于自定义Operator只是直接读取SQL文件的原始内容,没有触发Airflow的Jinja模板渲染机制,导致{{ parameters.table_name }}无法被替换为实际值。以下是修复方案:
1. 修改自定义Operator代码
更新custom_test_operator.py,加入模板渲染逻辑,将SQL文件中的模板变量替换为传入的参数:
from airflow.models.baseoperator import BaseOperator from airflow.utils.decorators import apply_defaults class MyTestOperator(BaseOperator): template_fields = ("my_file", "parameters",) @apply_defaults def __init__(self, my_file, parameters, *args, **kwargs): super(MyTestOperator, self).__init__(*args,**kwargs) self.my_file = my_file self.parameters = parameters def execute(self, context): print("目标文件:", self.my_file) # 读取SQL模板内容 with open(self.my_file, "r") as file: sql_template = file.read() # 使用Airflow的Jinja环境渲染模板,传入parameters变量 jinja_env = self.dag.get_template_env() rendered_sql = jinja_env.from_string(sql_template).render(parameters=self.parameters) print("渲染后的SQL语句:", rendered_sql) # 此处可添加实际执行SQL的逻辑(如使用DatabaseHook)
2. 可选优化:简化文件路径
由于你在DAG中设置了template_searchpath="/usr/local/airflow/include/sql/",可以将my_file的路径简化为文件名,避免硬编码全路径:
修改dag_custom_operator.py中的process_data任务:
process_data = MyTestOperator( task_id="process_data", my_file="query_1.sql", # 利用template_searchpath自动查找 parameters={ "table_name": "product" } )
验证效果
运行DAG后,execute方法会输出渲染后的SQL:
渲染后的SQL语句: SELECT COUNT(*) FROM product;
内容的提问来源于stack exchange,提问作者QDex
相关产品推荐
相关产品推荐

