You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.15 18:22:11