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

如何向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:

  1. 先安装依赖(Databricks笔记本中可通过魔法命令安装):
%pip install deltalake pandas
  1. 修改后的代码:
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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 21:50:25