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

Snowflake中动态扁平化未知结构JSON至表的实现方案

动态扁平化JSON并增量加载至Snowflake表的解决方案

针对无嵌套JSON的动态扁平化、自动新增列及增量加载需求,可通过Snowflake内置函数+存储过程自动化实现,无需预设列名,适配上游新增字段的场景。


核心思路

  1. 动态提取所有JSON键:通过OBJECT_KEYS展开JSON的所有键,聚合去重得到全量列名
  2. 同步目标表结构:对比目标表现有列,自动执行ALTER TABLE新增缺失列
  3. 增量加载数据:用控制表跟踪上次加载时间,仅处理新增数据,缺失字段自动填充NULL

具体实现步骤

1. 准备基础表

源数据表(存储Snowpipe Streaming导入的JSON)

CREATE OR REPLACE TABLE RAW_JSON_TABLE (
    RECORD_CONTENT VARIANT,
    LOAD_TIMESTAMP TIMESTAMP DEFAULT CURRENT_TIMESTAMP() -- 用于增量判断,也可替换为Snowpipe的OFFSET等标识
);

加载控制表(跟踪增量加载时间)

CREATE OR REPLACE TABLE LOAD_CONTROL_TABLE (
    TABLE_NAME VARCHAR(100) PRIMARY KEY,
    LAST_LOAD_TIME TIMESTAMP
);

2. 自动化处理存储过程

创建存储过程完成列同步+增量加载的全流程:

CREATE OR REPLACE PROCEDURE LOAD_FLATTENED_JSON()
RETURNS VARCHAR
LANGUAGE JAVASCRIPT
AS
$$
// 1. 获取本次新增数据中的所有JSON键
var getNewKeysSql = `
    SELECT DISTINCT KEY
    FROM RAW_JSON_TABLE,
         LATERAL FLATTEN(INPUT => OBJECT_KEYS(RECORD_CONTENT))
    WHERE LOAD_TIMESTAMP > COALESCE(
        (SELECT LAST_LOAD_TIME FROM LOAD_CONTROL_TABLE WHERE TABLE_NAME = 'FLATTENED_TARGET_TABLE'),
        '1970-01-01'::TIMESTAMP
    )
`;
var stmt = snowflake.createStatement({sqlText: getNewKeysSql});
var result = stmt.execute();
var newKeys = [];
while (result.next()) {
    newKeys.push(result.getColumnValue(1));
}

// 2. 检查目标表是否存在,不存在则创建
var checkTableExistsSql = `
    SELECT COUNT(*)
    FROM INFORMATION_SCHEMA.TABLES
    WHERE TABLE_NAME = 'FLATTENED_TARGET_TABLE'
      AND TABLE_SCHEMA = '${snowflake.getSchemaName()}'
      AND TABLE_CATALOG = '${snowflake.getDatabaseName()}'
`;
stmt = snowflake.createStatement({sqlText: checkTableExistsSql});
result = stmt.execute();
result.next();
var tableExists = result.getColumnValue(1) > 0;

if (!tableExists) {
    // 第一次运行,用所有键创建目标表
    var createTableSql = `CREATE OR REPLACE TABLE FLATTENED_TARGET_TABLE (${newKeys.map(key => `"${key}" VARIANT`).join(', ')})`;
    snowflake.createStatement({sqlText: createTableSql}).execute();
} else {
    // 目标表已存在,获取现有列并对比新增键,生成ALTER语句
    var getTargetColumnsSql = `
        SELECT COLUMN_NAME
        FROM INFORMATION_SCHEMA.COLUMNS
        WHERE TABLE_NAME = 'FLATTENED_TARGET_TABLE'
          AND TABLE_SCHEMA = '${snowflake.getSchemaName()}'
          AND TABLE_CATALOG = '${snowflake.getDatabaseName()}'
    `;
    stmt = snowflake.createStatement({sqlText: getTargetColumnsSql});
    result = stmt.execute();
    var targetColumns = new Set();
    while (result.next()) {
        targetColumns.add(result.getColumnValue(1));
    }

    var alterStatements = [];
    for (var key of newKeys) {
        if (!targetColumns.has(key)) {
            alterStatements.push(`ALTER TABLE FLATTENED_TARGET_TABLE ADD COLUMN "${key}" VARIANT;`);
        }
    }

    // 执行ALTER语句新增列
    for (var alterStmt of alterStatements) {
        snowflake.createStatement({sqlText: alterStmt}).execute();
    }
}

// 3. 获取上次加载时间,生成增量INSERT语句
var getLastLoadTimeSql = `
    SELECT COALESCE(MAX(LAST_LOAD_TIME), '1970-01-01'::TIMESTAMP)
    FROM LOAD_CONTROL_TABLE
    WHERE TABLE_NAME = 'FLATTENED_TARGET_TABLE'
`;
stmt = snowflake.createStatement({sqlText: getLastLoadTimeSql});
result = stmt.execute();
result.next();
var lastLoadTime = result.getColumnValue(1);

// 获取目标表所有列(含新增列)
var getAllColumnsSql = `
    SELECT COLUMN_NAME
    FROM INFORMATION_SCHEMA.COLUMNS
    WHERE TABLE_NAME = 'FLATTENED_TARGET_TABLE'
      AND TABLE_SCHEMA = '${snowflake.getSchemaName()}'
      AND TABLE_CATALOG = '${snowflake.getDatabaseName()}'
`;
stmt = snowflake.createStatement({sqlText: getAllColumnsSql});
result = stmt.execute();
var allColumns = [];
while (result.next()) {
    allColumns.push(result.getColumnValue(1));
}

// 生成INSERT的列和值映射
var columnsList = allColumns.map(col => `"${col}"`).join(', ');
var valuesList = allColumns.map(col => `RECORD_CONTENT:"${col}"`).join(', ');

var insertSql = `
    INSERT INTO FLATTENED_TARGET_TABLE (${columnsList})
    SELECT ${valuesList}
    FROM RAW_JSON_TABLE
    WHERE LOAD_TIMESTAMP > '${lastLoadTime}'
`;
snowflake.createStatement({sqlText: insertSql}).execute();

// 4. 更新控制表的最后加载时间
var updateControlSql = `
    MERGE INTO LOAD_CONTROL_TABLE t
    USING (SELECT 'FLATTENED_TARGET_TABLE' AS TABLE_NAME, CURRENT_TIMESTAMP() AS LAST_LOAD_TIME) s
    ON t.TABLE_NAME = s.TABLE_NAME
    WHEN MATCHED THEN UPDATE SET t.LAST_LOAD_TIME = s.LAST_LOAD_TIME
    WHEN NOT MATCHED THEN INSERT (TABLE_NAME, LAST_LOAD_TIME) VALUES (s.TABLE_NAME, s.LAST_LOAD_TIME)
`;
snowflake.createStatement({sqlText: updateControlSql}).execute();

return `处理完成:新增${newKeys.length}个字段,加载${insertSql.getRowCount()}条数据`;
$$;

3. 执行加载

每次有新增数据时,调用存储过程即可自动完成列同步和增量加载:

CALL LOAD_FLATTENED_JSON();

关键说明

  • 列类型选择VARIANT:兼容JSON所有数据类型(字符串、数字、布尔等),若后续需固定类型,可单独修改列类型
  • 增量判断:示例用LOAD_TIMESTAMP,可根据Snowpipe Streaming的实际场景替换为OFFSET或其他增量标识
  • 性能优化:仅提取新增数据中的JSON键,避免全表扫描,适合海量数据场景
  • 列名兼容:自动用双引号包裹JSON键,支持含特殊字符的键名

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 17:22:50