如何将70GB大JSON文件快速导入Snowflake表?
70GB大型JSON文件导入Snowflake方案
一、拆分大型JSON文件(基于ijson)
由于99.99999999%的数据集中在body数组,优先拆分该数组为多个小文件,实现并行加载。以下是Python示例代码:
import ijson import json import os # 创建存储拆分文件的目录 os.makedirs("body_chunks", exist_ok=True) # 读取原始大文件并提取元数据和body条目 with open("big.json", "r") as f: parser = ijson.parse(f) metadata = {} body_count = 0 for prefix, event, value in parser: # 提取顶层元数据 if prefix in ["entity_name", "entity_type", "last_updated_on", "version", "model_name", "model_id_type", "model_id", "model_market_type", "provider_references"]: metadata[prefix] = value # 提取每个body条目并写入单独文件 elif prefix == "body.item" and event == "end_map": # 每个文件包含元数据+单个body条目(保留原始结构) chunk_data = {**metadata, "body": [value]} with open(f"body_chunks/chunk_{body_count}.json", "w") as out_f: json.dump(chunk_data, out_f) body_count += 1 # 可选:每10000条合并为一个文件,减少文件数量 # if body_count % 10000 == 0: # print(f"Generated {body_count} chunks")
此代码会生成大量小JSON文件,每个文件包含完整的顶层元数据和单个body条目,适合Snowflake并行加载。
二、Snowflake表结构设计
根据数据嵌套关系,提供两种设计方案:
方案1:半结构化表(快速加载,适合灵活查询)
直接用VARIANT类型存储嵌套数据,无需提前拆分结构:
CREATE OR REPLACE TABLE ENTERPRISE_SEMI_STRUCTURED ( ENTITY_NAME STRING, ENTITY_TYPE STRING, LAST_UPDATED_ON DATE, VERSION STRING, MODEL_NAME STRING, MODEL_ID_TYPE STRING, MODEL_ID STRING, MODEL_MARKET_TYPE STRING, PROVIDER_REFERENCES VARIANT, BODY VARIANT );
方案2:规范化表(优化查询性能,适合复杂分析)
将嵌套结构拆分为多个关联表,避免冗余:
-- 顶层元数据表(仅1行) CREATE OR REPLACE TABLE ENTITY_METADATA ( ENTITY_ID INT DEFAULT 1, ENTITY_NAME STRING, ENTITY_TYPE STRING, LAST_UPDATED_ON DATE, VERSION STRING, MODEL_NAME STRING, MODEL_ID_TYPE STRING, MODEL_ID STRING, MODEL_MARKET_TYPE STRING, PROVIDER_REFERENCES VARIANT, PRIMARY KEY (ENTITY_ID) ); -- Body条目表 CREATE OR REPLACE TABLE BODY_ENTRIES ( BODY_ID INT AUTOINCREMENT, ENTITY_ID INT, NEGOTIATION_ARRANGEMENT STRING, NAME STRING, BILLING_CODE_TYPE STRING, BILLING_CODE_TYPE_VERSION STRING, BILLING_CODE STRING, DESCRIPTION STRING, PRIMARY KEY (BODY_ID), FOREIGN KEY (ENTITY_ID) REFERENCES ENTITY_METADATA(ENTITY_ID) ); -- 协商费率表 CREATE OR REPLACE TABLE NEGOTIATED_RATES ( NEGOTIATED_RATE_ID INT AUTOINCREMENT, BODY_ID INT, PRIMARY KEY (NEGOTIATED_RATE_ID), FOREIGN KEY (BODY_ID) REFERENCES BODY_ENTRIES(BODY_ID) ); -- 供应商组表 CREATE OR REPLACE TABLE PROVIDER_GROUPS ( PROVIDER_GROUP_ID INT AUTOINCREMENT, NEGOTIATED_RATE_ID INT, NPI ARRAY(INT), TIN_TYPE STRING, TIN_VALUE STRING, PRIMARY KEY (PROVIDER_GROUP_ID), FOREIGN KEY (NEGOTIATED_RATE_ID) REFERENCES NEGOTIATED_RATES(NEGOTIATED_RATE_ID) ); -- 协商价格表 CREATE OR REPLACE TABLE NEGOTIATED_PRICES ( NEGOTIATED_PRICE_ID INT AUTOINCREMENT, NEGOTIATED_RATE_ID INT, NEGOTIATED_TYPE STRING, NEGOTIATED_RATE FLOAT, EXPIRATION_DATE DATE, SERVICE_CODE ARRAY(STRING), BILLING_CLASS STRING, PRIMARY KEY (NEGOTIATED_PRICE_ID), FOREIGN KEY (NEGOTIATED_RATE_ID) REFERENCES NEGOTIATED_RATES(NEGOTIATED_RATE_ID) );
三、并行加载到Snowflake
1. 上传拆分文件到Snowflake Stage
首先创建内部Stage:
CREATE OR REPLACE STAGE BODY_CHUNKS_STAGE FILE_FORMAT = (TYPE = JSON);
使用SnowSQL上传文件:
PUT file:///path/to/body_chunks/*.json @BODY_CHUNKS_STAGE;
2. 加载到半结构化表
直接批量导入所有拆分文件:
COPY INTO ENTERPRISE_SEMI_STRUCTURED FROM @BODY_CHUNKS_STAGE FILE_FORMAT = (TYPE = JSON) MATCH_BY_COLUMN_NAME = CASE_INSENSITIVE;
3. 加载到规范化表
第一步:导入元数据
INSERT INTO ENTITY_METADATA (ENTITY_NAME, ENTITY_TYPE, LAST_UPDATED_ON, VERSION, MODEL_NAME, MODEL_ID_TYPE, MODEL_ID, MODEL_MARKET_TYPE, PROVIDER_REFERENCES) SELECT $1:entity_name::STRING, $1:entity_type::STRING, $1:last_updated_on::DATE, $1:version::STRING, $1:model_name::STRING, $1:model_id_type::STRING, $1:model_id::STRING, $1:model_market_type::STRING, $1:provider_references::VARIANT FROM @BODY_CHUNKS_STAGE/chunk_0.json;
第二步:导入Body条目
INSERT INTO BODY_ENTRIES (ENTITY_ID, NEGOTIATION_ARRANGEMENT, NAME, BILLING_CODE_TYPE, BILLING_CODE_TYPE_VERSION, BILLING_CODE, DESCRIPTION) SELECT 1 AS ENTITY_ID, body_item.value:negotiation_arrangement::STRING, body_item.value:name::STRING, body_item.value:billing_code_type::STRING, body_item.value:billing_code_type_version::STRING, body_item.value:billing_code::STRING, body_item.value:description::STRING FROM @BODY_CHUNKS_STAGE, LATERAL FLATTEN(input => $1:body) AS body_item;
第三步:导入嵌套数据
-- 导入协商费率 INSERT INTO NEGOTIATED_RATES (BODY_ID) SELECT be.BODY_ID FROM @BODY_CHUNKS_STAGE, LATERAL FLATTEN(input => $1:body) AS body_item, LATERAL FLATTEN(input => body_item.value:negotiated_rates) AS rate_item, JOIN BODY_ENTRIES be ON be.BILLING_CODE = body_item.value:billing_code::STRING AND be.NEGOTIATION_ARRANGEMENT = body_item.value:negotiation_arrangement::STRING; -- 导入供应商组 INSERT INTO PROVIDER_GROUPS (NEGOTIATED_RATE_ID, NPI, TIN_TYPE, TIN_VALUE) SELECT nr.NEGOTIATED_RATE_ID, pg_item.value:npi::ARRAY(INT), pg_item.value:tin.type::STRING, pg_item.value:tin.value::STRING FROM @BODY_CHUNKS_STAGE, LATERAL FLATTEN(input => $1:body) AS body_item, LATERAL FLATTEN(input => body_item.value:negotiated_rates) AS rate_item, LATERAL FLATTEN(input => rate_item.value:provider_groups) AS pg_item, JOIN BODY_ENTRIES be ON be.BILLING_CODE = body_item.value:billing_code::STRING, JOIN NEGOTIATED_RATES nr ON nr.BODY_ID = be.BODY_ID; -- 导入协商价格 INSERT INTO NEGOTIATED_PRICES (NEGOTIATED_RATE_ID, NEGOTIATED_TYPE, NEGOTIATED_RATE, EXPIRATION_DATE, SERVICE_CODE, BILLING_CLASS) SELECT nr.NEGOTIATED_RATE_ID, np_item.value:negotiated_type::STRING, np_item.value:negotiated_rate::FLOAT, np_item.value:expiration_date::DATE, np_item.value:service_code::ARRAY(STRING), np_item.value:billing_class::STRING FROM @BODY_CHUNKS_STAGE, LATERAL FLATTEN(input => $1:body) AS body_item, LATERAL FLATTEN(input => body_item.value:negotiated_rates) AS rate_item, LATERAL FLATTEN(input => rate_item.value:negotiated_prices) AS np_item, JOIN BODY_ENTRIES be ON be.BILLING_CODE = body_item.value:billing_code::STRING, JOIN NEGOTIATED_RATES nr ON nr.BODY_ID = be.BODY_ID;
四、加载优化建议
- 仓库规模:加载期间使用较大的仓库(如XL或XXL),加快并行处理速度。
- 文件大小:拆分后的文件保持在100MB-1GB之间,平衡并行度和文件管理成本。
- 自动 ingest:如果使用外部存储(如S3),开启自动 ingest 功能,文件上传后自动触发加载。
内容的提问来源于stack exchange,提问作者Alex
相关产品推荐
相关产品推荐

