如何通过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_
相关产品推荐
相关产品推荐

