如何用Airflow实现Snowflake Stage文件检测传感器(无S3权限)
实现Snowflake Stage文件检测的Airflow传感器
核心思路
通过Snowflake的SQL查询检查Stage内指定文件是否存在,结合Airflow的传感器组件执行查询并根据结果触发后续任务——无需直接访问S3,只要Snowflake Stage有权限读取S3文件即可。
方法一:用SqlSensor快速实现
1. 配置Snowflake连接
先在Airflow UI的「Admin → Connections」里添加Snowflake连接:Conn Type选Snowflake,填写账户、仓库、数据库、Schema、角色、认证信息(比如密钥或用户名密码),记好连接ID(比如snowflake_default)。
2. 编写检测SQL
Snowflake提供两种方式查询Stage文件:
- 方式A:用LIST命令(推荐,支持路径过滤)
SELECT COUNT(*) FROM TABLE(LIST(@你的Stage名称, '/目标文件所在路径/')) WHERE "name" = '要检测的文件名.csv';
如果要匹配一类文件,可改用通配符:WHERE "name" LIKE 'sales_2024%.csv'
- 方式B:查询INFORMATION_SCHEMA视图
SELECT COUNT(*) FROM INFORMATION_SCHEMA.STAGE_FILES WHERE STAGE_NAME = '你的Stage名称(大写)' AND RELATIVE_PATH = '目标文件相对路径/文件名.csv';
3. 在DAG中配置SqlSensor
直接用Airflow内置的SqlSensor执行上述SQL,当查询结果大于0时触发后续任务:
from airflow import DAG from airflow.sensors.sql import SqlSensor from datetime import datetime default_args = { 'owner': 'airflow', 'start_date': datetime(2024, 1, 1), } with DAG('snowflake_stage_file_check', default_args=default_args, schedule_interval='@daily', catchup=False) as dag: check_file_exists = SqlSensor( task_id='check_file_in_stage', conn_id='snowflake_default', # 替换为你的Snowflake连接ID sql=""" SELECT COUNT(*) FROM TABLE(LIST(@my_s3_stage, '/data/')) WHERE "name" = 'sales_data.csv'; """, poke_interval=300, # 每5分钟检查一次 timeout=3600, # 超时1小时后放弃 mode='poke', success=lambda records: int(records[0][0]) > 0, # 计数>0则文件存在 ) # 后续处理任务示例 process_snowflake_data = ... check_file_exists >> process_snowflake_data
方法二:封装自定义传感器(复用场景)
如果需要在多个DAG里复用检测逻辑,可以写个自定义传感器:
from airflow.sensors.base import BaseSensorOperator from airflow.providers.snowflake.hooks.snowflake import SnowflakeHook class SnowflakeStageSensor(BaseSensorOperator): def __init__( self, snowflake_conn_id: str, stage_name: str, file_path: str, poke_interval: int = 60, timeout: int = 3600, **kwargs, ): super().__init__(poke_interval=poke_interval, timeout=timeout, **kwargs) self.snowflake_conn_id = snowflake_conn_id self.stage_name = stage_name self.file_path = file_path def poke(self, context): hook = SnowflakeHook(snowflake_conn_id=self.snowflake_conn_id) sql = f""" SELECT COUNT(*) FROM TABLE(LIST(@{self.stage_name})) WHERE "name" = '{self.file_path}'; """ records = hook.get_records(sql) return int(records[0][0]) > 0 # DAG中使用自定义传感器 with DAG(...) as dag: check_stage_file = SnowflakeStageSensor( task_id='check_snowflake_stage_file', snowflake_conn_id='snowflake_default', stage_name='my_s3_stage', file_path='data/sales_data.csv', poke_interval=300, )
注意事项
- 确保Snowflake角色拥有Stage的
USAGE权限,以及查询INFORMATION_SCHEMA的权限(如果用方式B)。 - 避免过于频繁的轮询,防止给Snowflake造成不必要的负载。
- 如果Stage是外部S3 Stage,只要Snowflake本身能访问S3,这个逻辑就能正常工作——不需要Airflow拥有S3权限。
内容的提问来源于stack exchange,提问作者user3735871
相关产品推荐
相关产品推荐

