如何使用Snowflake存储过程重组异构数据表的列顺序
实现方案
核心逻辑
- 先给所有行加有序行号,保证原始数据顺序不丢失
- 识别空行(所有列值均为NULL的行)作为数据段分隔符,给每行分配所属数据段ID
- 提取每个数据段的第一行作为该段的字段映射表头
- 按统一目标字段规则,将各段数据映射对齐,缺失字段补NULL
实现步骤
1. 预处理数据(添加行号与段ID)
首先创建带元数据行号的中间表,确保行顺序和导入时一致:
CREATE OR REPLACE TEMP TABLE countries_pp_with_seg AS WITH all_rows AS ( SELECT METADATA$FILE_ROW_NUMBER AS row_num, -- 导入时保留的原始行号,顺序完全匹配原始文件 column1, column2, column3, column4, column5 FROM countries_pp ), -- 标记空行 empty_row_flag AS ( SELECT *, CASE WHEN column1 IS NULL AND column2 IS NULL AND column3 IS NULL AND column4 IS NULL AND column5 IS NULL THEN 1 ELSE 0 END AS is_empty_row FROM all_rows ), -- 计算段ID,每遇到空行段ID+1 seg_calculate AS ( SELECT *, SUM(is_empty_row) OVER (ORDER BY row_num ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS seg_id FROM empty_row_flag WHERE is_empty_row = 0 -- 过滤掉空行 ) SELECT * FROM seg_calculate;
2. 提取各段表头映射
CREATE OR REPLACE TEMP TABLE seg_header_map AS -- 每个段的第一行是表头 WITH seg_headers AS ( SELECT seg_id, column1 AS c1_name, column2 AS c2_name, column3 AS c3_name, column4 AS c4_name, column5 AS c5_name, MIN(row_num) OVER (PARTITION BY seg_id) AS header_row_num FROM countries_pp_with_seg ) SELECT seg_id, OBJECT_CONSTRUCT( c1_name, 1, c2_name, 2, c3_name, 3, c4_name, 4, c5_name, 5 ) AS col_pos_map -- 存储当前段字段名对应的原始列位置 FROM seg_headers WHERE row_num = header_row_num;
3. 数据对齐映射生成最终结果
CREATE OR REPLACE TABLE final_aligned_result AS WITH seg_data_rows AS ( -- 过滤掉每个段的表头行,只保留数据行 SELECT a.*, b.col_pos_map FROM countries_pp_with_seg a JOIN seg_header_map b ON a.seg_id = b.seg_id WHERE a.row_num > (SELECT MIN(row_num) FROM countries_pp_with_seg WHERE seg_id = a.seg_id) ) SELECT -- 按目标列顺序映射,不存在的字段返回NULL CASE WHEN col_pos_map:Year IS NOT NULL THEN GET(a, 'COLUMN'||col_pos_map:Year::INT) END AS Year, CASE WHEN col_pos_map:Country IS NOT NULL THEN GET(a, 'COLUMN'||col_pos_map:Country::INT) END AS Country, CASE WHEN col_pos_map:City IS NOT NULL THEN GET(a, 'COLUMN'||col_pos_map:City::INT) END AS City, CASE WHEN col_pos_map:Name IS NOT NULL THEN GET(a, 'COLUMN'||col_pos_map:Name::INT) END AS Name, CASE WHEN col_pos_map:Age IS NOT NULL THEN GET(a, 'COLUMN'||col_pos_map:Age::INT) END AS Age FROM seg_data_rows a;
4. 封装为存储过程
把上述逻辑封装成存储过程,调用后直接生成最终对齐表:
CREATE OR REPLACE PROCEDURE proc_columns_matching() RETURNS STRING LANGUAGE JAVASCRIPT AS $$ // 执行预处理SQL const sql_steps = [ `CREATE OR REPLACE TEMP TABLE countries_pp_with_seg AS WITH all_rows AS ( SELECT METADATA$FILE_ROW_NUMBER AS row_num, column1, column2, column3, column4, column5 FROM countries_pp ), empty_row_flag AS ( SELECT *, CASE WHEN column1 IS NULL AND column2 IS NULL AND column3 IS NULL AND column4 IS NULL AND column5 IS NULL THEN 1 ELSE 0 END AS is_empty_row FROM all_rows ), seg_calculate AS ( SELECT *, SUM(is_empty_row) OVER (ORDER BY row_num ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS seg_id FROM empty_row_flag WHERE is_empty_row = 0 ) SELECT * FROM seg_calculate`, `CREATE OR REPLACE TEMP TABLE seg_header_map AS WITH seg_headers AS ( SELECT seg_id, column1 AS c1_name, column2 AS c2_name, column3 AS c3_name, column4 AS c4_name, column5 AS c5_name, MIN(row_num) OVER (PARTITION BY seg_id) AS header_row_num FROM countries_pp_with_seg ) SELECT seg_id, OBJECT_CONSTRUCT( c1_name, 1, c2_name, 2, c3_name, 3, c4_name, 4, c5_name, 5 ) AS col_pos_map FROM seg_headers WHERE row_num = header_row_num`, `CREATE OR REPLACE TABLE final_aligned_result AS WITH seg_data_rows AS ( SELECT a.*, b.col_pos_map FROM countries_pp_with_seg a JOIN seg_header_map b ON a.seg_id = b.seg_id WHERE a.row_num > (SELECT MIN(row_num) FROM countries_pp_with_seg WHERE seg_id = a.seg_id) ) SELECT CASE WHEN col_pos_map:Year IS NOT NULL THEN GET(a, 'COLUMN'||col_pos_map:Year::INT) END AS Year, CASE WHEN col_pos_map:Country IS NOT NULL THEN GET(a, 'COLUMN'||col_pos_map:Country::INT) END AS Country, CASE WHEN col_pos_map:City IS NOT NULL THEN GET(a, 'COLUMN'||col_pos_map:City::INT) END AS City, CASE WHEN col_pos_map:Name IS NOT NULL THEN GET(a, 'COLUMN'||col_pos_map:Name::INT) END AS Name, CASE WHEN col_pos_map:Age IS NOT NULL THEN GET(a, 'COLUMN'||col_pos_map:Age::INT) END AS Age FROM seg_data_rows a` ]; for (let sql of sql_steps) { snowflake.createStatement({sqlText: sql}).execute(); } return "处理完成,结果已写入表 final_aligned_result"; $$;
调用方法
CALL proc_columns_matching(); -- 查看结果 SELECT * FROM final_aligned_result;
注意事项
- 如果你的原始表没有导入时保留的
METADATA$FILE_ROW_NUMBER,可以用ROW_NUMBER() OVER(ORDER BY 1)替代,但要确保Snowflake返回的行顺序和原始文件一致,建议导入时显式添加行号字段更稳妥。 - 目标列如果有新增,只需要在最终SELECT部分新增对应字段的映射逻辑即可。
内容的提问来源于stack exchange,提问作者Catapultaa
相关产品推荐
相关产品推荐

