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

如何将本地拆分JSON的Python代码适配到Databricks平台?

问题分析

你尝试的Databricks代码存在核心问题:df_json是Spark DataFrame对象,并非Python字典,因此不能直接调用items()方法进行切片操作,这也是后续无法完成写入的根源。以下提供两种适配Databricks场景的解决方案,分别对应单文件原生处理和分布式Spark处理两种需求。


方案一:沿用Python原生逻辑(适合单个大型JSON文件)

Databricks支持通过/dbfs/前缀直接访问DBFS存储,因此可以复用本地的字典处理逻辑,仅需调整文件路径:

import json
import itertools

# 读取DBFS中的JSON文件,路径需添加/dbfs/前缀
with open('/dbfs/mnt/BigData_JSONFiles/SampleDatafilefrombigfile.json', 'r') as fp:
    data = json.loads(fp.read())

# 拆分头部(前8个键值对)和详情内容
d1 = dict(itertools.islice(data.items(), 8))
d2 = dict(itertools.islice(data.items(), 8, len(data)))

# 写入拆分后的文件到DBFS
with open("/dbfs/mnt/BigData_JSONFiles/new_test_header.json", "w") as header_file:
    json.dump(d1, header_file, indent=2)
with open("/dbfs/mnt/BigData_JSONFiles/new_test_detail.json", "w") as detail_file:
    json.dump(d2, detail_file, indent=2)

说明

  • 逻辑与本地代码完全一致,仅修改路径适配DBFS规则
  • 适合单个大型JSON文件,但需注意文件大小不能超过单节点内存限制

方案二:用PySpark DataFrame处理(适合分布式场景)

如果要处理超大型文件或批量JSON文件,推荐使用Spark分布式能力,通过DataFrame API实现拆分:

from pyspark.sql import functions as F

# 读取多线格式的JSON文件
df_json = spark.read.option("multiline", "true").json("/mnt/BigData_JSONFiles/SampleDatafilefrombigfile.json")

# 定义头部列(对应原JSON前8个字段)
header_cols = [
    "reporting_entity_name", "reporting_entity_type", "plan_name",
    "plan_id_type", "plan_id", "plan_market_type", "last_updated_on", "version"
]
# 自动识别详情列(排除头部列)
detail_cols = [col for col in df_json.columns if col not in header_cols]

# 拆分头部和详情DataFrame
df_header = df_json.select(*header_cols)
df_detail = df_json.select(*detail_cols)

# 写入头部JSON(coalesce(1)用于合并为单个文件,可选)
df_header.coalesce(1) \
    .write \
    .option("multiline", "true") \
    .mode("overwrite") \
    .json("/mnt/BigData_JSONFiles/header_output")

# 写入详情JSON
df_detail.coalesce(1) \
    .write \
    .option("multiline", "true") \
    .mode("overwrite") \
    .json("/mnt/BigData_JSONFiles/detail_output")

说明

  • coalesce(1)会将输出合并为单个文件,默认Spark会生成多个分区文件,可根据需求选择是否保留
  • mode("overwrite")表示覆盖已有文件,也可替换为append(追加)或ignore(忽略)模式
  • Spark会将文件写入指定目录,若需要指定文件名,可后续用dbutils.fs.mv命令重命名生成的分区文件

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 05:30:50