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

如何用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 18:43:32