如何将Spark DataFrame转换为含指定数组与列表的嵌套字典
如何将Spark DataFrame转换为指定结构的嵌套字典?
需求说明
需要把给定的Spark DataFrame转换成包含accepted数组的嵌套字典结构,每个数组元素包含issuer、Recipient嵌套字典,以及additional_fields子数组。
示例DataFrame
columns = ["name","reason","cgc","limit","email","address","message","type","value"] data = [("Paulo", "La Fava","123456","0","p@p.com.br","avenue A","msg txt 1","string","low"), ("Pedro", "Petrus","123457","20.00","pop@petrus.com.br","avenue A","msg txt 2","string", "average"), ("Saulo", "Salix","123458","150.00","python@salix.com.br","avenue B","msg txt 3","string","high")] df = spark.createDataFrame(data).toDF(*columns)
方法一:用Spark内置函数构造结构(适合大数据量)
利用Spark的struct和array函数在DataFrame层面直接构建嵌套结构,再转成JSON后解析为字典,这种方法适合分布式环境下的大数据处理。
from pyspark.sql import functions as F import json # 构建嵌套结构的DataFrame nested_df = df.select( F.array( F.struct( F.struct("name", "reason", "cgc").alias("issuer"), F.struct("limit", "email", "address").alias("Recipient"), F.array(F.struct("message", "type", "value")).alias("additional_fields") ) ).alias("accepted") ) # 将DataFrame转为JSON字符串,再解析为字典 result_dict = json.loads(nested_df.toJSON().first()) print(json.dumps(result_dict, indent=2))
方法二:Python手动构造(适合小数据量)
如果数据量不大,可以先将DataFrame的行数据收集到本地,再通过Python循环手动构造目标字典结构,灵活性更高。
import json # 收集DataFrame所有行到本地 rows = df.collect() # 初始化结果字典 result = {"accepted": []} # 遍历每行构造嵌套结构 for row in rows: item = { "issuer": { "name": row.name, "reason": row.reason, "cgc": row.cgc }, "Recipient": { "limit": row.limit, "email": row.email, "address": row.address }, "additional_fields": [ { "message": row.message, "type": row.type, "value": row.value } ] } result["accepted"].append(item) # 打印格式化后的结果 print(json.dumps(result, indent=2))
输出结果
两种方法都能生成符合要求的嵌套字典(以完整数据为例):
{ "accepted": [ { "issuer": { "name": "Paulo", "reason": "La Fava", "cgc": "123456" }, "Recipient": { "limit": "0", "email": "p@p.com.br", "address": "avenue A" }, "additional_fields": [ { "message": "msg txt 1", "type": "string", "value": "low" } ] }, { "issuer": { "name": "Pedro", "reason": "Petrus", "cgc": "123457" }, "Recipient": { "limit": "20.00", "email": "pop@petrus.com.br", "address": "avenue A" }, "additional_fields": [ { "message": "msg txt 2", "type": "string", "value": "average" } ] }, { "issuer": { "name": "Saulo", "reason": "Salix", "cgc": "123458" }, "Recipient": { "limit": "150.00", "email": "python@salix.com.br", "address": "avenue B" }, "additional_fields": [ { "message": "msg txt 3", "type": "string", "value": "high" } ] } ] }
内容的提问来源于stack exchange,提问作者SrKartcheski
相关产品推荐
相关产品推荐

