如何从嵌套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
相关产品推荐
相关产品推荐

