Snowflake任务设置on_error=continue时,如何获取文件拷贝错误通知?
解决思路与方案
方案1:用存储过程封装COPY,主动检查错误并发通知
既然ON_ERROR=CONTINUE会让任务整体标记为成功,那就自己在逻辑里主动检查COPY的错误情况,触发SNS通知,同时保证任务正常完成不终止。
具体操作:
- 确认已配置好AWS SNS的通知集成(示例中集成名为
MY_SNS_INTEGRATION)。 - 编写存储过程封装COPY逻辑:
- 执行带
ON_ERROR=CONTINUE的COPY命令,确保单个文件失败不中断任务。 - 查询
COPY_HISTORY视图,获取本次任务触发的COPY错误统计。 - 若错误数大于0,调用
SYSTEM$SEND_NOTIFICATION主动发送邮件通知。
- 执行带
示例代码:
CREATE OR REPLACE PROCEDURE COPY_WITH_ERROR_ALERT() RETURNS VARCHAR LANGUAGE SQL AS $$ DECLARE copy_sql VARCHAR := 'COPY INTO my_target_table FROM @my_s3_stage FILE_FORMAT = (FORMAT_NAME = my_file_format) ON_ERROR = CONTINUE'; current_task_id VARCHAR := CURRENT_TASK(); error_num INTEGER; BEGIN -- 执行COPY命令,单个文件失败不终止 EXECUTE IMMEDIATE copy_sql; -- 查询最近5分钟内当前任务触发的COPY错误数 SELECT COUNT(*) INTO error_num FROM TABLE(INFORMATION_SCHEMA.COPY_HISTORY(TABLE_NAME => 'my_target_table', START_TIME => DATEADD('minute', -5, CURRENT_TIMESTAMP()))) WHERE TASK_NAME = current_task_id AND LOAD_STATUS = 'LOADED_WITH_ERRORS'; -- 存在错误则发送SNS通知 IF error_num > 0 THEN CALL SYSTEM$SEND_NOTIFICATION( 'MY_SNS_INTEGRATION', 'Snowflake COPY任务有文件加载失败', '任务ID: ' || current_task_id || '\n失败文件数: ' || error_num || '\n可查看COPY_HISTORY获取详情' ); END IF; RETURN 'COPY执行完毕,共' || error_num || '个文件加载失败'; END; $$; -- 修改原有任务,改为调用该存储过程 ALTER TASK my_copy_task SET WAREHOUSE = my_warehouse SCHEDULE = 'USING CRON 0 * * * * UTC' AS CALL COPY_WITH_ERROR_ALERT();
方案2:用错误日志表+监控任务
给COPY命令指定错误日志表,将所有失败文件的信息持久化存储,再创建独立任务监控该表的新增记录,有错误时触发通知。
步骤:
- 创建错误日志表:
CREATE OR REPLACE TABLE my_copy_error_log ( file_name VARCHAR, error_msg VARCHAR, error_code INTEGER, load_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP() );
- 修改COPY命令,将错误信息写入日志表:
COPY INTO my_target_table FROM @my_s3_stage FILE_FORMAT = (FORMAT_NAME = my_file_format) ON_ERROR = CONTINUE ERROR_LOG_TABLE = my_copy_error_log;
- 创建流监控日志表的新增数据:
CREATE OR REPLACE STREAM my_copy_error_stream ON TABLE my_copy_error_log APPEND_ONLY = TRUE;
- 创建定时任务,检测流中新增错误并发送通知:
CREATE OR REPLACE TASK my_error_alert_task WAREHOUSE = my_warehouse SCHEDULE = 'USING CRON 5 * * * * UTC' -- 比COPY任务晚5分钟执行,确保COPY完成 WHEN SYSTEM$STREAM_HAS_DATA('my_copy_error_stream') AS BEGIN -- 临时存储新增错误记录 CREATE OR REPLACE TEMP TABLE temp_errors AS SELECT file_name, error_msg FROM my_copy_error_stream; -- 发送错误通知 CALL SYSTEM$SEND_NOTIFICATION( 'MY_SNS_INTEGRATION', 'Snowflake COPY任务检测到加载错误', '错误详情:\n' || LISTAGG('文件: ' || file_name || ' | 错误: ' || error_msg, '\n') WITHIN GROUP (ORDER BY load_time) ); -- 归档错误记录,避免重复通知 INSERT INTO my_copy_error_archive SELECT * FROM temp_errors; END;
方案3:单独创建任务监控COPY历史
搭建独立的定时任务,定期查询最近的COPY执行记录,若发现状态为LOADED_WITH_ERRORS的任务,立即发送通知。
示例代码:
CREATE OR REPLACE TASK my_copy_monitor_task WAREHOUSE = my_warehouse SCHEDULE = 'USING CRON 10 * * * * UTC' AS BEGIN -- 抓取最近1小时内存在错误的COPY记录 CREATE OR REPLACE TEMP TABLE recent_failures AS SELECT TASK_NAME, file_name, error_msg, start_time FROM TABLE(INFORMATION_SCHEMA.COPY_HISTORY(TABLE_NAME => 'my_target_table', START_TIME => DATEADD('hour', -1, CURRENT_TIMESTAMP()))) WHERE LOAD_STATUS = 'LOADED_WITH_ERRORS'; -- 存在错误则触发通知 IF (SELECT COUNT(*) FROM recent_failures) > 0 THEN CALL SYSTEM$SEND_NOTIFICATION( 'MY_SNS_INTEGRATION', 'Snowflake COPY任务存在加载失败', '任务名: ' || (SELECT DISTINCT TASK_NAME FROM recent_failures) || '\n最近失败时间: ' || (SELECT MAX(start_time) FROM recent_failures) || '\n失败文件数: ' || (SELECT COUNT(*) FROM recent_failures) ); END IF; END;
关键注意事项
- 确保
SYSTEM$SEND_NOTIFICATION函数权限正常,且SNS集成已正确关联AWS SNS主题,Snowflake角色具备对应操作权限。 - 根据COPY任务的执行频率调整监控任务的调度时间,避免遗漏或重复通知。
- 查询
COPY_HISTORY时,合理设置时间范围,避免查询到无关的历史记录。
内容的提问来源于stack exchange,提问作者Estrobelai
相关产品推荐
相关产品推荐

