如何在PySpark中将多表关联视图的数据转换为指定嵌套JSON结构
PySpark构造嵌套JSON结构实现方案
前置依赖导入
from pyspark.sql import functions as F
单region场景实现
假设你关联生成的视图名为store_sale_view,已包含所有需要的原始字段,直接按目标结构嵌套构造即可:
# 读取关联后的视图数据 df = spark.table("store_sale_view") # 构造目标嵌套结构 nested_df = df.select( # 最外层store字段 "store", # region数组字段 F.array( F.struct( F.col("region_id").alias("id"), # 嵌套location结构 F.struct("lat", "long").alias("location"), # 嵌套country结构 F.struct(F.col("country_name").alias("name")).alias("country"), # 嵌套city结构 F.struct( F.col("city_name").alias("name"), F.col("city_region").alias("region") ).alias("city"), F.col("year").alias("year "), "nivel", "slug" ) ).alias("region"), # 嵌套sale结构 F.struct(F.col("sale_store").alias("store")).alias("sale") )
同门店多region场景实现
如果同一个门店对应多条region数据,需要先分组再聚合为数组:
nested_df = df.groupBy("store", "sale_store").agg( F.collect_list( F.struct( F.col("region_id").alias("id"), F.struct("lat", "long").alias("location"), F.struct(F.col("country_name").alias("name")).alias("country"), F.struct( F.col("city_name").alias("name"), F.col("city_region").alias("region") ).alias("city"), F.col("year").alias("year "), "nivel", "slug" ) ).alias("region") ).select( "store", "region", F.struct(F.col("sale_store").alias("store")).alias("sale") )
输出JSON文件
# 每行输出一个JSON对象,适合后续大数据场景读取 nested_df.write.mode("overwrite").json("你的存储路径/output") # 若需要输出为单个数组格式的JSON文件,可增加参数调整 nested_df.coalesce(1).write.option("lineSep", ",").mode("overwrite").json("你的存储路径/single_output")
注意事项
- 请根据你实际视图的字段名替换代码中的字段名,保持和你业务表的字段对应即可
- 字段别名需要和你目标JSON的key完全一致,包括空格、大小写
内容的提问来源于stack exchange,提问作者Molly
相关产品推荐
相关产品推荐

