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

Airflow任务未因Snowflake存储过程运行时错误失败的问题排查

解决Snowflake存储过程报错时Airflow任务不失败的问题

问题根源

当前Airflow任务未捕获到Snowflake存储过程的错误,核心原因是存储过程内部处理了异常但未重新抛出,导致cursor.execute()执行后没有向外传递错误;其次是手动管理连接时未主动检查执行结果。

修复步骤

1. 调整Snowflake存储过程,确保错误能抛出

修改存储过程代码,在异常捕获块中主动抛出错误,让外部调用者感知到失败:

CREATE OR REPLACE PROCEDURE SP_LOAD_DATA_S3_TO_SNOWFLAKE(application_name VARCHAR, config_db VARCHAR, config_schema VARCHAR)
RETURNS VARCHAR
LANGUAGE JAVASCRIPT
AS
$$
try {
    // 原有数据加载逻辑
    // ...
    return 'SUCCESS';
} catch (err) {
    // 新增:将错误重新抛出,让Airflow捕获
    throw `Stored procedure failed: ${err.message}`;
}
$$;

2. 修改Airflow Python任务代码

推荐使用SnowflakeHook的run方法代替手动管理连接和游标,该方法会自动捕获Snowflake执行异常:

def load():
    try:
        dwh_hook = SnowflakeHook(
            snowflake_conn_id=Variable.get("SNOWFLAKE_CONN_ID"),
            warehouse=Variable.get("SNOWFLAKE_WAREHOUSE"),
            database=app_db,
            role=Variable.get("SNOWFLAKE_ROLE"),
            schema=app_stg_schema
        )
        
        # 构造存储过程调用语句(建议用参数绑定避免SQL注入)
        call_sp = f"call {db}.{schema}.SP_LOAD_DATA_S3_TO_SNOWFLAKE('{application_name}','{config_db}','{config_schema}')"
        
        # 使用run方法执行,自动处理异常和提交
        dwh_hook.run(call_sp, autocommit=True)
        
        move_files_within_bucket(file_names)
        print("file_names:", file_names)
    except Exception as e:
        print("An error occurred:", str(e))
        raise e

如果必须手动管理游标,需主动检查执行结果或错误状态:

def load():
    try:
        dwh_hook = SnowflakeHook(...)
        con = dwh_hook.get_conn()
        cursor = con.cursor()

        call_sp = f"call {db}.{schema}.SP_LOAD_DATA_S3_TO_SNOWFLAKE('{application_name}','{config_db}','{config_schema}')"
        cursor.execute(call_sp)
        
        # 获取存储过程返回结果并检查
        result = cursor.fetchone()
        if result and not result[0].startswith('SUCCESS'):
            raise Exception(f"存储过程执行失败,返回结果:{result[0]}")
        
        # 额外检查游标错误状态
        if hasattr(cursor, '_error') and cursor._error:
            raise Exception(f"Snowflake游标错误:{cursor._error}")
        
        con.commit()
        cursor.close()
        con.close()
        
        move_files_within_bucket(file_names)
        print("file_names:", file_names)
    except Exception as e:
        print("An error occurred:", str(e))
        raise e

3. 额外注意事项

  • 避免手动拼接SQL语句,建议使用SnowflakeHook的参数绑定功能防止SQL注入
  • 确认Airflow的default_args中email_on_failure已开启,确保失败时收到通知
  • 检查Snowflake仓库、角色权限是否足够执行存储过程

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.07 12:35:25