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

Snowpipe如何加载S3中Base64编码的双层JSON数据到Snowflake表

解决方案

核心问题原因

你之前的查询逻辑存在两处错误导致解析失败:

  1. 给Stage指定了JSON类型的文件格式后,Snowflake会自动将文件内容解析为JSON结构化对象,此时$1不再是完整的原始文件字符串,无法直接做Base64解码
  2. 外层文件本身就是明文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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 06:27:04