Snowpipe如何加载S3中Base64编码的双层JSON数据到Snowflake表
解决方案
核心问题原因
你之前的查询逻辑存在两处错误导致解析失败:
- 给Stage指定了JSON类型的文件格式后,Snowflake会自动将文件内容解析为JSON结构化对象,此时
$1不再是完整的原始文件字符串,无法直接做Base64解码 - 外层文件本身就是明文JSON格式,不需要对整个文件内容做Base64解码,仅内层的
payload字段值才是Base64编码的内容
操作步骤
1. 创建原始文本读取格式
创建一个纯文本文件格式,用于将整个S3文件内容作为单个原始字符串列读取,避免Snowflake提前解析JSON结构:
CREATE OR REPLACE FILE FORMAT RAW_SINGLE_TEXT_FORMAT TYPE = 'CSV' FIELD_DELIMITER = NONE RECORD_DELIMITER = NONE ESCAPE = NONE ESCAPE_UNENCLOSED_FIELD = NONE QUOTE = NONE;
2. 修正COPY INTO语句
调整解析逻辑,先解析外层JSON,再对payload字段单独做Base64解码和内层JSON解析:
COPY INTO TEST_TABLE ( event_key1, event_key2, event_key3, metadata1, -- 可选:如果需要存入context中的元数据可保留 metadata2 ) FROM ( SELECT -- 先解析外层JSON -> 取payload字段 -> Base64解码 -> 解析内层JSON -> 取对应字段 PARSE_JSON(BASE64_DECODE_STRING(PARSE_JSON($1):payload)):event_key1::STRING AS event_key1, PARSE_JSON(BASE64_DECODE_STRING(PARSE_JSON($1):payload)):event_key2::STRING AS event_key2, PARSE_JSON(BASE64_DECODE_STRING(PARSE_JSON($1):payload)):event_key3::STRING AS event_key3, -- 可选:读取context中的元数据 PARSE_JSON($1):context.metadata1::STRING AS metadata1, PARSE_JSON($1):context.metadata2::STRING AS metadata2 FROM @TEST_STAGE ) FILE_FORMAT = (FORMAT_NAME = 'RAW_SINGLE_TEXT_FORMAT');
3. 创建自动接入的Snowpipe
基于上面的COPY逻辑创建开启自动摄入的Snowpipe,之后在AWS S3控制台配置事件通知,将新增文件事件推送到Snowpipe返回的SQS队列即可实现自动接入:
CREATE OR REPLACE PIPE TEST_S3_PIPE AUTO_INGEST = TRUE AS COPY INTO TEST_TABLE ( event_key1, event_key2, event_key3 ) FROM ( SELECT PARSE_JSON(BASE64_DECODE_STRING(PARSE_JSON($1):payload)):event_key1::STRING AS event_key1, PARSE_JSON(BASE64_DECODE_STRING(PARSE_JSON($1):payload)):event_key2::STRING AS event_key2, PARSE_JSON(BASE64_DECODE_STRING(PARSE_JSON($1):payload)):event_key3::STRING AS event_key3 FROM @TEST_STAGE ) FILE_FORMAT = (FORMAT_NAME = 'RAW_SINGLE_TEXT_FORMAT');
验证操作
执行以下语句确认配置正确:
-- 查看Snowpipe的SQS队列ARN,用于配置S3事件通知 SHOW PIPES LIKE 'TEST_S3_PIPE'; -- 查看管道最近的加载历史 SELECT * FROM TABLE(INFORMATION_SCHEMA.PIPE_USAGE_HISTORY( DATE_RANGE_START=>DATEADD('day',-1,CURRENT_DATE()), PIPE_NAME=>'TEST_S3_PIPE' ));
内容的提问来源于stack exchange,提问作者maverickz
相关产品推荐
相关产品推荐

