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

如何通过Airflow Jinja与Params实现SnowflakeOperator动态S3 Key?

解决SnowflakeOperator动态生成S3 Key的方案

方法1:用PythonOperator预处理参数并通过XCom传递

先通过PythonOperator手动渲染包含execution_date的S3 Key模板,将结果存入XCom,再在SnowflakeOperator中拉取使用:

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.providers.snowflake.operators.snowflake import SnowflakeOperator
from datetime import datetime
from jinja2 import Template

def process_s3_key(**context):
    # 定义带模板变量的S3 Key格式
    s3_key_template = "foo/blah/{{ execution_date }}/test_file.csv"
    # 渲染模板,指定日期格式
    rendered_s3_key = Template(s3_key_template).render(
        execution_date=context['execution_date'].strftime('%Y-%m-%d')
    )
    # 将渲染结果推送到XCom
    context['ti'].xcom_push(key='rendered_s3_key', value=rendered_s3_key)

with DAG(
    dag_id='snowflake_unload_dynamic_key',
    start_date=datetime(2023, 8, 21),
    schedule_interval=None,
) as dag:
    process_key = PythonOperator(
        task_id='process_s3_key',
        python_callable=process_s3_key,
        provide_context=True,
    )

    snowflake_unload = SnowflakeOperator(
        task_id='snowflake_unload',
        snowflake_conn_id='your_snowflake_connection_id',
        sql='path/to/your/sql_template.sql',
        params={'s3_key': '{{ ti.xcom_pull(task_ids="process_s3_key", key="rendered_s3_key") }}'},
    )

    process_key >> snowflake_unload

方法2:自定义SnowflakeOperator子类自动渲染params

继承SnowflakeOperator,重写execute方法,在执行SQL前自动渲染params中的模板变量:

from airflow.providers.snowflake.operators.snowflake import SnowflakeOperator
from jinja2 import Template
from airflow.utils.context import Context

class RenderedParamsSnowflakeOperator(SnowflakeOperator):
    def execute(self, context: Context):
        # 遍历params,渲染所有字符串类型的模板变量
        rendered_params = {}
        for key, value in self.params.items():
            if isinstance(value, str):
                rendered_params[key] = Template(value).render(**context)
            else:
                rendered_params[key] = value
        self.params = rendered_params
        # 调用父类的execute方法执行SQL
        return super().execute(context)

使用自定义Operator时,直接在params中传入带模板的S3 Key即可:

from airflow import DAG
from datetime import datetime

with DAG(
    dag_id='snowflake_unload_custom_operator',
    start_date=datetime(2023, 8, 21),
    schedule_interval=None,
) as dag:
    snowflake_unload = RenderedParamsSnowflakeOperator(
        task_id='snowflake_unload',
        snowflake_conn_id='your_snowflake_connection_id',
        sql='path/to/your/sql_template.sql',
        params={'s3_key': 'foo/blah/{{ execution_date.strftime("%Y-%m-%d") }}/test_file.csv'},
    )

方法3:手动渲染完整SQL语句

直接用Jinja2加载SQL模板,先处理S3 Key参数,再生成最终SQL传给SnowflakeOperator:

from airflow import DAG
from airflow.providers.snowflake.operators.snowflake import SnowflakeOperator
from jinja2 import Environment, FileSystemLoader
from datetime import datetime

with DAG(
    dag_id='snowflake_unload_render_sql',
    start_date=datetime(2023, 8, 21),
    schedule_interval=None,
) as dag:
    def get_rendered_sql(**context):
        # 加载SQL模板文件
        env = Environment(loader=FileSystemLoader('path/to/sql/templates/directory'))
        sql_template = env.get_template('your_sql_template.sql')
        
        # 渲染S3 Key参数
        s3_key = Template('foo/blah/{{ execution_date }}/test_file.csv').render(
            execution_date=context['execution_date'].strftime('%Y-%m-%d')
        )
        
        # 渲染完整SQL语句
        return sql_template.render(params={'s3_key': s3_key})

    snowflake_unload = SnowflakeOperator(
        task_id='snowflake_unload',
        snowflake_conn_id='your_snowflake_connection_id',
        sql=get_rendered_sql,
        provide_context=True,
    )

内容的提问来源于stack exchange,提问作者mad_

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 19:17:39