如何向Delta Table追加数据?求改造现有JSON转CSV存储代码
将JSON.gz数据写入Delta Table的实现方案
方案一:使用PySpark(推荐,适配Databricks分布式环境)
原代码直接加载整个JSON文件到内存的方式,在处理大体积的json.gz文件时容易出现内存溢出问题,PySpark的分布式处理能力更适合这类场景,同时原生支持Delta Lake:
# 导入PySpark相关模块 from pyspark.sql import SparkSession # 初始化SparkSession(Databricks环境中可省略,默认已配置) spark = SparkSession.builder \ .appName("JSONtoDelta") \ .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \ .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \ .getOrCreate() # 读取JSON.gz文件,PySpark自动识别压缩格式 df = spark.read.json("/dbfs/mnt/costtransparency/HealthPlans/UHC/2022-11/2022-11-01_United-HealthCare-Services--Inc-_Third-Party-Administrator_Winstead_CSP-903-C746_in-network-rates.json.gz") # 提取'in_network'数组字段并展开为单独行,对应原代码遍历数组的逻辑 expanded_df = df.selectExpr("explode(in_network) as in_network_data").select("in_network_data.*") # 写入Delta Table,指定存储路径 expanded_df.write \ .format("delta") \ .mode("overwrite") # 按需选择模式:append(追加)/ignore(忽略重复)/error(存在则报错) .save("/dbfs/mnt/delta_tables/uhc_in_network_rates") # 可选:将Delta表注册到Metastore,方便后续用SQL查询 spark.sql("CREATE TABLE IF NOT EXISTS uhc_in_network_rates USING DELTA LOCATION '/dbfs/mnt/delta_tables/uhc_in_network_rates'")
关键说明:
explode函数负责将数组类型的in_network拆分为单行数据,和原代码遍历数组的逻辑等价- Databricks环境默认集成Delta Lake,无需额外安装依赖
方案二:使用Pandas + Delta-rs(适合小文件场景)
如果数据量较小,可基于原代码的Pandas逻辑修改,通过delta-rs库写入Delta Table:
- 先安装依赖(Databricks笔记本中可通过魔法命令安装):
%pip install deltalake pandas
- 修改后的代码:
import json import gzip import pandas as pd from deltalake import write_deltalake # 修正原代码的压缩文件读取逻辑,用gzip.open处理.gz格式 with gzip.open('/dbfs/mnt/costtransparency/HealthPlans/UHC/2022-11/2022-11-01_United-HealthCare-Services--Inc-_Third-Party-Administrator_Winstead_CSP-903-C746_in-network-rates.json.gz', 'rt') as f: d = json.load(f) employee_data = d['in_network'] # 转换为Pandas DataFrame df = pd.DataFrame(employee_data) # 写入Delta Table write_deltalake( '/dbfs/mnt/delta_tables/uhc_in_network_rates_pandas', df, mode='overwrite' # 按需选择模式 )
注意事项:
- 原代码直接用
open读取.json.gz会报错,必须用gzip.open并指定rt(文本模式)读取 - 此方案仅适合小数据集,大数据量下会出现内存压力
内容的提问来源于stack exchange,提问作者Shrikant Shejwal_IND
相关产品推荐
相关产品推荐

