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

Snowflake任务设置on_error=continue时,如何获取文件拷贝错误通知?

解决思路与方案

方案1:用存储过程封装COPY,主动检查错误并发通知

既然ON_ERROR=CONTINUE会让任务整体标记为成功,那就自己在逻辑里主动检查COPY的错误情况,触发SNS通知,同时保证任务正常完成不终止。

具体操作:

  1. 确认已配置好AWS SNS的通知集成(示例中集成名为MY_SNS_INTEGRATION)。
  2. 编写存储过程封装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命令指定错误日志表,将所有失败文件的信息持久化存储,再创建独立任务监控该表的新增记录,有错误时触发通知。

步骤:

  1. 创建错误日志表:
CREATE OR REPLACE TABLE my_copy_error_log (
  file_name VARCHAR,
  error_msg VARCHAR,
  error_code INTEGER,
  load_time TIMESTAMP DEFAULT CURRENT_TIMESTAMP()
);
  1. 修改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;
  1. 创建流监控日志表的新增数据:
CREATE OR REPLACE STREAM my_copy_error_stream
ON TABLE my_copy_error_log
APPEND_ONLY = TRUE;
  1. 创建定时任务,检测流中新增错误并发送通知:
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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 09:10:57