如何让Airflow的SQLExecuteQueryOperator执行Snowflake PUT时用本地文件系统而非S3 Accelerate
解决Apache Airflow中SQLExecuteQueryOperator执行Snowflake PUT命令时强制使用本地文件系统传输的问题
问题核心是Airflow的Snowflake钩子默认采用S3作为文件传输代理,导致PUT命令触发S3 Accelerate调用。要改用本地文件系统传输,可通过以下两种方式配置:
方法1:修改Snowflake连接全局配置
在Airflow的Snowflake连接(snowflake_default)的Extra字段中添加JSON配置:
{"file_transfer_type": "LOCAL"}
此配置会让所有使用该连接的任务默认采用本地文件系统传输。
方法2:在任务实例中单独指定参数
若不想修改全局连接,可直接在SQLExecuteQueryOperator中通过hook_params传递配置,修改后的完整DAG代码如下:
from airflow import DAG from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator import pendulum put_command = """ PUT file://filepath/test.csv @my_snowflake_stage SOURCE_COMPRESSION=AUTO_DETECT; """ default_args = { 'owner': 'airflow', 'depends_on_past': False, 'email_on_failure': False, 'email_on_retry': False, } with DAG( dag_id="example_put_to_snowflake", default_args=default_args, start_date=pendulum.yesterday("UTC"), schedule_interval=None, catchup=False, ) as dag: put_task = SQLExecuteQueryOperator( task_id="put_file_to_stage", conn_id="snowflake_default", sql=put_command, # 指定文件传输类型为本地 hook_params={"file_transfer_type": "LOCAL"} )
关键说明
file_transfer_type参数用于控制Snowflake钩子的传输代理类型,可选值包括S3(默认)、LOCAL、AZURE等;- 设置为
LOCAL后,PUT命令会直接利用Airflow Worker节点的本地文件系统完成上传,不再走S3路径; - 需确保Worker节点能访问
file://filepath/test.csv指定的本地文件路径,否则会触发文件不存在的错误。
内容的提问来源于stack exchange,提问作者asasisekar
相关产品推荐
相关产品推荐

