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

Snowflake存储过程无法处理S3文件,请求协助排查

Snowflake存储过程文件存在性检查及加载故障排查

问题背景

需要创建Snowflake存储过程,实现从S3加载当日最新数据到表中,同时添加了以下错误处理逻辑:

  • 文件不存在时写入日志表
  • 文件存在但加载失败时标记日志
  • 加载成功时记录日志
    但调用存储过程时出现错误,核心问题是即使文件实际存在,也无法正确获取文件计数,导致后续逻辑异常。

怀疑有问题的代码片段

file_count_1 := (
SELECT COUNT(*)
FROM TABLE(
FUNCTION('LIST', '@@DB_NAME.SCHEMA_NAME.S3_SF_DEV/Test/FILE/DATA', { PATTERN => dynamic_file_pattern })
)

完整存储过程代码

CREATE OR REPLACE PROCEDURE COPY_DATA_TO_RAW()
RETURNS VARCHAR
LANGUAGE SQL
EXECUTE AS OWNER
AS
$$
DECLARE
  dynamic_file_pattern VARCHAR;
  dynamic_file_pattern_1 VARCHAR;
  copy_command VARCHAR;
  copy_command_1 VARCHAR;
  file_count_1 INTEGER;
  file_count_2 INTEGER;
BEGIN

  -- Construct the dynamic file patterns
  dynamic_file_pattern := '.*EQ_TEM_TST' || TO_VARCHAR(CURRENT_DATE(), 'YYYYMMDD') || 'T.*[.]csv';
  dynamic_file_pattern_1 := '.*EQ_TST' || TO_VARCHAR(CURRENT_DATE(), 'YYYYMMDD') || 'T.*[.]csv';

  -- Check if files exist using LIST command
  file_count_1 := (
    SELECT COUNT(*)
    FROM TABLE(
      FUNCTION('LIST', '@@DB_NAME.SCHEMA_NAME.S3_SF_DEV/Test/FILE/DATA', { PATTERN => dynamic_file_pattern })
    )
  );

  file_count_2 := (
    SELECT COUNT(*)
    FROM TABLE(
      FUNCTION('LIST', '@@DB_NAME.SCHEMA_NAME.S3_SF_DEV/Test/FILE/DATA', { PATTERN => dynamic_file_pattern_1 })
    )
  );


  -- If any file is not found, log the error and stop execution
  IF (file_count_1 = 0 OR file_count_2 = 0) THEN
    INSERT INTO AUDIT_TABLE (TABLE_NAME, STATUS, DATE_MODIFIED)
    VALUES ('TABLE1', 'Latest files not found.', CURRENT_TIMESTAMP);
    RETURN 'Latest files not found.';
  END IF;

  -- Truncate tables before loading data
  TRUNCATE TABLE TABLE1;
  TRUNCATE TABLE TABLE2;

  -- Construct the COPY INTO command dynamically for TABLE1
  copy_command := '
    COPY INTO TABLE1 (
      INTERNAL, NAME, DESCRIPTION, START, ACTION, CLASS, VISIBILITY, ATTRIBUTES, ACATTRIBUTES
    )
    FROM (
      SELECT
        t.$1 AS INTERNAL,
        t.$2 AS NAME,
        t.$3 AS DESCRIPTION,
        t.$4 AS START,
        CASE
          WHEN t.$5 = ''Create'' THEN 1
          WHEN t.$5 = ''Modify'' THEN 2
          WHEN t.$5 = ''Delete'' THEN 3
        END AS ACTION,
        t.$6 AS CLASS,
        t.$7 AS VISIBILITY,
        t.$8 AS ATTRIBUTES,
        t.$9 AS ACATTRIBUTES
      FROM @@DB_NAME.SCHEMA_NAME.S3_SF_DEV/Test/FILE/DATA
      (PATTERN => '''' || dynamic_file_pattern || '''') AS t
    )
    FILE_FORMAT = (
      TYPE = ''CSV'',
      COMPRESSION = ''NONE'',
      SKIP_HEADER = 1,
      FIELD_DELIMITER = ''|'',
      FIELD_OPTIONALLY_ENCLOSED_BY = ''"''
    )
  ';
  -- Construct the COPY INTO command dynamically for TABLE2
  copy_command_1 := '
    COPY INTO TABLE2 (
      INTERNAL, NAME, DESC, ORG, EQUTN, ISEQUA, ACTIVE, FILTER, ATTRIBUTES, ACATTRIBUTES
    )
    FROM (
      SELECT
        t.$1 AS INTERNAL,
        t.$2 AS NAME,
        IFF(t.$3 = '''', ''<NULL>'', t.$3) AS DESC,
        t.$4 AS ORG,
        t.$5 AS EQUTN,
        CASE WHEN t.$6 = ''YES'' THEN 1 ELSE 0 END AS ISEQUA,
        CASE WHEN t.$7 = ''YES'' THEN 1 ELSE 0 END AS ACTIVE,
        t.$8 AS FILTER,
        t.$9 AS ATTRIBUTES,
        t.$10 AS ACATTRIBUTES
      FROM @@DB_NAME.SCHEMA_NAME.S3_SF_DEV/Test/FILE/DATA
      (PATTERN => '''' || dynamic_file_pattern_1 || '''') AS t
    )
    FILE_FORMAT = (
      TYPE = ''CSV'',
      COMPRESSION = ''NONE'',
      SKIP_HEADER = 1,
      FIELD_DELIMITER = ''|'',
      FIELD_OPTIONALLY_ENCLOSED_BY = ''"''
    )
  ';

  -- Try loading data, and handle any failures
  BEGIN
    EXECUTE IMMEDIATE copy_command;
    EXECUTE IMMEDIATE copy_command_1;

    INSERT INTO AUDIT_TABLE (TABLE_NAME, STATUS, DATE_MODIFIED)
    VALUES ('TABLE1', 'SUCCESS', CURRENT_TIMESTAMP);
  EXCEPTION WHEN OTHER THEN
    INSERT INTO AUDIT_TABLE (TABLE_NAME, STATUS, DATE_MODIFIED)
    VALUES ('TABLE1', 'LOAD FAILED', CURRENT_TIMESTAMP);
    RETURN 'LOAD FAILED';
  END;

  RETURN 'LOAD COMPLETED SUCCESSFULLY';
END;
$$;

问题排查与修复方案

1. 文件计数查询的语法错误

SQL存储过程中调用LIST函数时,无法直接通过FUNCTION()语法传递动态变量,Snowflake无法解析这种写法,导致计数始终为0。需改用动态SQL执行处理模式匹配参数:

-- 修复file_count_1的动态查询
EXECUTE IMMEDIATE 'SELECT COUNT(*) FROM TABLE(LIST(''@@DB_NAME.SCHEMA_NAME.S3_SF_DEV/Test/FILE/DATA'', PATTERN => ''' || dynamic_file_pattern || '''))' INTO file_count_1;

-- 修复file_count_2的动态查询
EXECUTE IMMEDIATE 'SELECT COUNT(*) FROM TABLE(LIST(''@@DB_NAME.SCHEMA_NAME.S3_SF_DEV/Test/FILE/DATA'', PATTERN => ''' || dynamic_file_pattern_1 || '''))' INTO file_count_2;

2. COPY命令中的字符串拼接错误

原代码中PATTERN参数的引号转义逻辑错误,导致生成的SQL语句语法无效,同时FIELD_OPTIONALLY_ENCLOSED_BY的HTML转义符"在动态SQL中需替换为直接转义的引号:

-- 修复后的TABLE1 COPY命令片段
copy_command := '
    ...
      FROM @@DB_NAME.SCHEMA_NAME.S3_SF_DEV/Test/FILE/DATA
      (PATTERN => ''' || dynamic_file_pattern || ''') AS t
    ...
    FILE_FORMAT = (
      ...
      FIELD_OPTIONALLY_ENCLOSED_BY = ''"''
    )
  ';

3. 错误处理的完整性优化

原代码仅记录TABLE1的日志状态,未覆盖TABLE2,建议调整日志插入逻辑,分别记录两个表的状态:

-- 成功时记录两个表的状态
INSERT INTO AUDIT_TABLE (TABLE_NAME, STATUS, DATE_MODIFIED)
VALUES ('TABLE1', 'SUCCESS', CURRENT_TIMESTAMP),
       ('TABLE2', 'SUCCESS', CURRENT_TIMESTAMP);

-- 异常时记录两个表的失败状态
EXCEPTION WHEN OTHER THEN
    INSERT INTO AUDIT_TABLE (TABLE_NAME, STATUS, DATE_MODIFIED)
    VALUES ('TABLE1', 'LOAD FAILED', CURRENT_TIMESTAMP),
           ('TABLE2', 'LOAD FAILED', CURRENT_TIMESTAMP);
    RETURN 'LOAD FAILED';

4. 变量替换注意事项

确保@@DB_NAME.SCHEMA_NAME替换为实际的数据库和模式名称,避免使用未解析的占位符。

修复后的完整存储过程

CREATE OR REPLACE PROCEDURE COPY_DATA_TO_RAW()
RETURNS VARCHAR
LANGUAGE SQL
EXECUTE AS OWNER
AS
$$
DECLARE
  dynamic_file_pattern VARCHAR;
  dynamic_file_pattern_1 VARCHAR;
  copy_command VARCHAR;
  copy_command_1 VARCHAR;
  file_count_1 INTEGER;
  file_count_2 INTEGER;
BEGIN

  -- 构造动态文件匹配模式
  dynamic_file_pattern := '.*EQ_TEM_TST' || TO_VARCHAR(CURRENT_DATE(), 'YYYYMMDD') || 'T.*[.]csv';
  dynamic_file_pattern_1 := '.*EQ_TST' || TO_VARCHAR(CURRENT_DATE(), 'YYYYMMDD') || 'T.*[.]csv';

  -- 动态检查文件存在性并计数
  EXECUTE IMMEDIATE 'SELECT COUNT(*) FROM TABLE(LIST(''@@DB_NAME.SCHEMA_NAME.S3_SF_DEV/Test/FILE/DATA'', PATTERN => ''' || dynamic_file_pattern || '''))' INTO file_count_1;
  EXECUTE IMMEDIATE 'SELECT COUNT(*) FROM TABLE(LIST(''@@DB_NAME.SCHEMA_NAME.S3_SF_DEV/Test/FILE/DATA'', PATTERN => ''' || dynamic_file_pattern_1 || '''))' INTO file_count_2;

  -- 任一文件不存在则记录日志并终止
  IF (file_count_1 = 0 OR file_count_2 = 0) THEN
    INSERT INTO AUDIT_TABLE (TABLE_NAME, STATUS, DATE_MODIFIED)
    VALUES ('TABLE1', 'Latest files not found.', CURRENT_TIMESTAMP),
           ('TABLE2', 'Latest files not found.', CURRENT_TIMESTAMP);
    RETURN 'Latest files not found.';
  END IF;

  -- 加载前清空表
  TRUNCATE TABLE TABLE1;
  TRUNCATE TABLE TABLE2;

  -- 构造TABLE1的COPY命令
  copy_command := '
    COPY INTO TABLE1 (
      INTERNAL, NAME, DESCRIPTION, START, ACTION, CLASS, VISIBILITY, ATTRIBUTES, ACATTRIBUTES
    )
    FROM (
      SELECT
        t.$1 AS INTERNAL,
        t.$2 AS NAME,
        t.$3 AS DESCRIPTION,
        t.$4 AS START,
        CASE
          WHEN t.$5 = ''Create'' THEN 1
          WHEN t.$5 = ''Modify'' THEN 2
          WHEN t.$5 = ''Delete'' THEN 3
        END AS ACTION,
        t.$6 AS CLASS,
        t.$7 AS VISIBILITY,
        t.$8 AS ATTRIBUTES,
        t.$9 AS ACATTRIBUTES
      FROM @@DB_NAME.SCHEMA_NAME.S3_SF_DEV/Test/FILE/DATA
      (PATTERN => ''' || dynamic_file_pattern || ''') AS t
    )
    FILE_FORMAT = (
      TYPE = ''CSV'',
      COMPRESSION = ''NONE'',
      SKIP_HEADER = 1,
      FIELD_DELIMITER = ''|'',
      FIELD_OPTIONALLY_ENCLOSED_BY = ''"''
    )
  ';

  -- 构造TABLE2的COPY命令
  copy_command_1 := '
    COPY INTO TABLE2 (
      INTERNAL, NAME, DESC, ORG, EQUTN, ISEQUA, ACTIVE, FILTER, ATTRIBUTES, ACATTRIBUTES
    )
    FROM (
      SELECT
        t.$1 AS INTERNAL,
        t.$2 AS NAME,
        IFF(t.$3 = '''', ''<NULL>'', t.$3) AS DESC,
        t.$4 AS ORG,
        t.$5 AS EQUTN,
        CASE WHEN t.$6 = ''YES'' THEN 1 ELSE 0 END AS ISEQUA,
        CASE WHEN t.$7 = ''YES'' THEN 1 ELSE 0 END AS ACTIVE,
        t.$8 AS FILTER,
        t.$9 AS ATTRIBUTES,
        t.$10 AS ACATTRIBUTES
      FROM @@DB_NAME.SCHEMA_NAME.S3_SF_DEV/Test/FILE/DATA
      (PATTERN => ''' || dynamic_file_pattern_1 || ''') AS t
    )
    FILE_FORMAT = (
      TYPE = ''CSV'',
      COMPRESSION = ''NONE'',
      SKIP_HEADER = 1,
      FIELD_DELIMITER = ''|'',
      FIELD_OPTIONALLY_ENCLOSED_BY = ''"''
    )
  ';

  -- 执行加载并处理异常
  BEGIN
    EXECUTE IMMEDIATE copy_command;
    EXECUTE IMMEDIATE copy_command_1;

    INSERT INTO AUDIT_TABLE (TABLE_NAME, STATUS, DATE_MODIFIED)
    VALUES ('TABLE1', 'SUCCESS', CURRENT_TIMESTAMP),
           ('TABLE2', 'SUCCESS', CURRENT_TIMESTAMP);
  EXCEPTION WHEN OTHER THEN
    INSERT INTO AUDIT_TABLE (TABLE_NAME, STATUS, DATE_MODIFIED)
    VALUES ('TABLE1', 'LOAD FAILED', CURRENT_TIMESTAMP),
           ('TABLE2', 'LOAD FAILED', CURRENT_TIMESTAMP);
    RETURN 'LOAD FAILED';
  END;

  RETURN 'LOAD COMPLETED SUCCESSFULLY';
END;
$$;

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 09:49:50