Spark Scala中分组后如何生成带id键的结构体数组列
实现方案
你当前写法的问题是直接把原始字符串类型的itemId传入聚合函数,得到的自然是字符串集合,只要先把itemId包装成要求的结构体结构,再做聚合即可,不需要自定义函数或者额外后置处理。
推荐写法:聚合阶段直接构造目标结构
用内置的结构体构造函数先把itemId转成{"id": 对应值}的结构,再传入集合收集函数:
注意:
collect_set会对同组内的itemId去重,如果你需要保留重复值,把collect_set替换为collect_list即可。
DataFrame API 写法
from pyspark.sql import functions as F db.select("customerIdMarketplace", "itemId") .groupBy("customerIdMarketplace") .agg( F.collect_set(F.struct(F.col("itemId").alias("id"))).alias("items") )
Spark SQL 写法
SELECT customerIdMarketplace, collect_set(named_struct('id', itemId)) AS items FROM 你的源表名 GROUP BY customerIdMarketplace
备选写法:聚合后二次转换结构
如果你已经跑完了原代码,拿到了items为字符串数组的中间结果,也可以用数组高阶函数直接遍历转换,不用重新跑分组聚合:
from pyspark.sql import functions as F # 假设old_df是你原代码输出的结果 old_df.withColumn( "items", F.expr("transform(items, item_val -> named_struct('id', item_val))") )
两种写法最终输出的结构和你给出的期望示例完全一致。
内容的提问来源于stack exchange,提问作者Jericho Sims
相关产品推荐
相关产品推荐

