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

如何从嵌套JSON生成DataFrame(不使用pandas处理大数据场景)

基于Spark的实现方案(完全无pandas依赖,适配大数据量场景)

你可以直接用Spark原生API解析该嵌套JSON结构,全程分布式运行,不会出现单节点内存瓶颈:

单/少量JSON文件场景实现

from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, IntegerType, StringType

# 初始化Spark会话
spark = SparkSession.builder.appName("NestedJsonToDf").getOrCreate()

# 读取原始JSON文件,multiLine参数适配多行格式化的JSON
raw_df = spark.read.json("your_json_file_path.json", multiLine=True)

# 提取表结构和数据
table_info = raw_df.select("tables.*").first()
# 动态提取列名和类型生成Schema,无需硬编码适配字段变更
type_mapping = {"Int": IntegerType(), "String": StringType()}
schema = StructType([
    StructField(col["name"], type_mapping[col["type"]], nullable=True)
    for col in table_info["columns"]
])
row_data = table_info["rows"]

# 生成目标DataFrame
result_df = spark.createDataFrame(row_data, schema=schema)

# 验证输出
result_df.show()

运行后输出和你要求的格式完全一致:

+----------+------------+--------------+
|EmployeeID|EmployeeName|DepartmentName|
+----------+------------+--------------+
|       123|    John Doe|            IT|
|       234|    Jane Doe|            HR|
+----------+------------+--------------+

大量同结构JSON文件场景优化

如果需要处理批量JSON文件,避免将全量数据拉取到Driver端,可以用RDD分布式解析:

import json
from pyspark.sql import SparkSession
from pyspark.sql.types import StructType, StructField, IntegerType, StringType

spark = SparkSession.builder.appName("BatchJsonParse").getOrCreate()

# 提前从样例文件提取Schema复用,避免重复解析
sample_json = json.load(open("sample_json_path.json"))
type_mapping = {"Int": IntegerType(), "String": StringType()}
schema = StructType([
    StructField(col["name"], type_mapping[col["type"]], nullable=True)
    for col in sample_json["tables"]["columns"]
])

# 分布式读取并解析所有JSON文件
def parse_single_json(line):
    try:
        json_data = json.loads(line)
        return json_data["tables"]["rows"]
    except Exception:
        return []

raw_rdd = spark.sparkContext.textFile("/your/json/directory/*.json")
rows_rdd = raw_rdd.flatMap(parse_single_json)

result_df = spark.createDataFrame(rows_rdd, schema=schema)

该方案所有解析逻辑都在Executor端执行,支持TB级数据量处理,无单节点内存压力。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 05:54:04