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
相关产品推荐
相关产品推荐

