在PySpark中基于多表关联视图生成含JSON数组的目标文档
PySpark 生成指定嵌套JSON结构实现方案
首先假设你已生成的关联视图名为product_view,结构如下:
| 字段名 | 类型 | 说明 |
|---|---|---|
| id | int | 主键id |
| color | string | 单个颜色值 |
| size | string | 单个尺寸值 |
| hash | int | 同id下全局唯一的哈希值 |
具体实现代码
from pyspark.sql.functions import col, struct, collect_list # 1. 读取已生成的视图为DataFrame df = spark.sql("SELECT id, color, size, hash FROM product_view") # 2. 构造嵌套结构体、按顶层字段分组聚合 result_df = df.groupBy("id", "hash") \ .agg( # 聚合color为 [{"color":"xxx"}] 结构 collect_list(struct(col("color").alias("color"))).alias("colors"), # 聚合size为 [{"size":"xxx"}] 结构 collect_list(struct(col("size").alias("size"))).alias("sizes") ) # 3. 输出为JSON格式 # 方案A:直接写入文件,每行对应一个JSON对象 result_df.write.mode("overwrite").json("你的输出路径") # 方案B:如果需要在Driver端拿到JSON字符串结果 json_result = result_df.toJSON().collect() for item in json_result: print(item)
常见适配调整
- 如果需要去重数组中的重复值,将
collect_list替换为collect_set即可 - 如果hash字段需要动态计算,可以在分组前使用
hash(col("自定义字段1"), col("自定义字段2"))自定义生成 - 如果视图中颜色和尺寸是已经聚合好的数组格式,可先使用
explode拆分为单行单个值后再执行上述聚合逻辑
内容的提问来源于stack exchange,提问作者Molly
相关产品推荐
相关产品推荐

