如何使用PySpark或Python将JSON列表拆分为单个JSON对象
问题描述
我有如下JSON列表:
[ { "files": 0, "data": [ {"name": "RFC", "value": "XXXXXXX", "attId": 1}, {"name": "NOMBRE", "value": "JOSE", "attId": 2}, {"name": "APELLIDO PATERNO", "value": "MONTIEL", "attId": 3}, {"name": "APELLIDO MATERNO", "value": "MENDOZA", "attId": 4}, {"name": "FECHA NACIMIENTO", "value": "1989-02-04", "attId": 5} ], "dirId": 1, "docId": 4, "structure": { "name": "personales", "folioId": 22 } }, { "files": 0, "data": [ {"name": "CALLE", "value": "AMOR", "attId": 6}, {"name": "No. EXTERIOR", "value": "4", "attId": 7}, {"name": "No. INTERIOR", "value": "2", "attId": 8}, {"name": "C.P.", "value": "55060", "attId": 9}, {"name": "ENTIDAD", "value": "ESTADO DE MEXICO", "attId": 10}, {"name": "MUNICIPIO", "value": "ECATEPEC", "attId": 11}, {"name": "COLONIA", "value": "INDUSTRIAL", "attId": 12} ], "dirId": 1, "docId": 4, "structure": { "name": "direccion", "folioId": 22 } } ]
需要将该JSON列表转换为独立的单个JSON对象并分别执行处理,请问如何使用PySpark或Python实现该需求?
一、Python原生实现
直接用Python内置的json模块解析列表,遍历每个对象完成处理,适合小数据量场景。
代码示例
import json # 示例JSON字符串,实际可从文件读取(用json.load()) json_str = ''' [ { "files": 0, "data": [ {"name": "RFC", "value": "XXXXXXX", "attId": 1}, {"name": "NOMBRE", "value": "JOSE", "attId": 2}, {"name": "APELLIDO PATERNO", "value": "MONTIEL", "attId": 3}, {"name": "APELLIDO MATERNO", "value": "MENDOZA", "attId": 4}, {"name": "FECHA NACIMIENTO", "value": "1989-02-04", "attId": 5} ], "dirId": 1, "docId": 4, "structure": { "name": "personales", "folioId": 22 } }, { "files": 0, "data": [ {"name": "CALLE", "value": "AMOR", "attId": 6}, {"name": "No. EXTERIOR", "value": "4", "attId": 7}, {"name": "No. INTERIOR", "value": "2", "attId": 8}, {"name": "C.P.", "value": "55060", "attId": 9}, {"name": "ENTIDAD", "value": "ESTADO DE MEXICO", "attId": 10}, {"name": "MUNICIPIO", "value": "ECATEPEC", "attId": 11}, {"name": "COLONIA", "value": "INDUSTRIAL", "attId": 12} ], "dirId": 1, "docId": 4, "structure": { "name": "direccion", "folioId": 22 } } ] ''' # 解析为Python列表,每个元素是独立的JSON对象(字典) json_objects = json.loads(json_str) # 逐个处理每个对象 for index, obj in enumerate(json_objects): print(f"=== 处理第 {index+1} 个对象 ===") # 示例处理1:提取结构名称 print(f"结构类型: {obj['structure']['name']}") # 示例处理2:将data数组转为键值对字典 data_dict = {item['name']: item['value'] for item in obj['data']} print("字段键值对:", data_dict) # 可添加自定义逻辑:比如写入文件、校验字段等
关键说明
- 用
json.loads()(字符串)或json.load()(文件)解析后,直接得到包含所有独立对象的列表。 - 遍历列表即可对每个对象单独操作,处理逻辑完全自定义。
二、PySpark实现
适合大数据量场景,通过Spark的DataFrame API解析列表、拆分对象并批量处理。
方法1:直接基于列表创建DataFrame
如果已经拿到Python格式的JSON列表,直接转为DataFrame,每行对应一个独立对象:
from pyspark.sql import SparkSession from pyspark.sql.functions import udf, col from pyspark.sql.types import MapType, StringType # 初始化Spark会话 spark = SparkSession.builder.appName("SplitJSONList").getOrCreate() # 原始JSON列表(同Python示例中的json_objects) json_objects = json.loads(json_str) # 创建DataFrame,每行是一个独立JSON对象 df = spark.createDataFrame(json_objects) # 自定义处理函数:将data数组转为键值对字典 def process_data(data_list): return {item['name']: item['value'] for item in data_list} # 注册UDF process_data_udf = udf(process_data, MapType(StringType(), StringType())) # 处理数据,新增data_map字段 df_processed = df.withColumn("data_map", process_data_udf(col("data"))) # 查看结果 df_processed.select("structure.name", "data_map").show(truncate=False)
方法2:解析JSON字符串列表
如果原始数据是包含整个JSON列表的字符串,先解析为数组再拆分:
from pyspark.sql.functions import from_json, explode from pyspark.sql.types import ArrayType, StructType, StructField, IntegerType, StringType # 定义单个JSON对象的Schema data_item_schema = StructType([ StructField("name", StringType()), StructField("value", StringType()), StructField("attId", IntegerType()) ]) structure_schema = StructType([ StructField("name", StringType()), StructField("folioId", IntegerType()) ]) json_obj_schema = StructType([ StructField("files", IntegerType()), StructField("data", ArrayType(data_item_schema)), StructField("dirId", IntegerType()), StructField("docId", IntegerType()), StructField("structure", structure_schema) ]) # 模拟原始数据:包含JSON列表字符串的DataFrame df_raw = spark.createDataFrame([(json_str,)], ["json_list_str"]) # 解析字符串为JSON数组 df_parsed = df_raw.withColumn("json_array", from_json(col("json_list_str"), ArrayType(json_obj_schema))) # 拆分数组,每个元素转为一行(独立对象) df_exploded = df_parsed.select(explode(col("json_array")).alias("single_obj")) # 提取字段处理 df_final = df_exploded.select( col("single_obj.structure.name").alias("structure_type"), col("single_obj.data") ) df_final.show(truncate=False)
关键说明
explode函数是核心:将数组类型的列拆分为多行,每行对应一个独立的JSON对象。- 可结合Spark内置函数或UDF完成批量处理,支持分布式计算,适合海量数据。
内容的提问来源于stack exchange,提问作者Harshith K R
相关产品推荐
相关产品推荐

