Spark 2.4.0版本下无需UDF实现Struct数组按subs值降序排序的方法
无需UDF实现Struct数组按subs降序排序
当然可以!在Spark 2.4.0版本中,我们完全不需要自定义UDF,用内置的函数组合就能轻松实现你要的效果。下面给你两种可行的实现方式:
方式一:窗口函数方案
这种方式利用窗口函数严格控制每个城市内部的排序范围,逻辑清晰可控:
from pyspark.sql import functions as F from pyspark.sql.window import Window # 定义窗口:按城市唯一标识分组,按subs字段降序排序 window_spec = Window.partitionBy("city_id", "city_name").orderBy(F.col("cel_item.subs").desc()) # 执行排序逻辑 sorted_df = df.withColumn("cel_item", F.explode("cel")) \ .withColumn("sorted_cel", F.collect_list("cel_item").over(window_spec)) \ .groupBy("city_id", "city_name") \ .agg(F.first("sorted_cel").alias("cel")) # 查看最终结果 sorted_df.show(2, False)
逻辑拆解
- 拆分数组:
F.explode("cel")把每个城市对应的cel数组拆分成单独行,每行对应一个(carr, subs)结构体,命名为cel_item。 - 窗口排序收集:窗口的
partitionBy确保我们只在同一个城市的范围内处理数据,orderBy指定按subs降序排列;再通过collect_list在窗口范围内收集排序后的结构体,生成每个城市的排序后数组。 - 恢复原结构:因为同一个城市的所有拆分行会生成完全相同的排序后数组,用
F.first提取唯一结果,恢复成一行对应一个城市的原始结构。
方式二:拆解-排序-聚合方案
这种方式更简洁直观,直接通过全局排序后聚合完成需求:
from pyspark.sql import functions as F sorted_df = df.select("city_id", "city_name", F.explode("cel").alias("cel_item")) \ .orderBy("city_id", "city_name", F.col("cel_item.subs").desc()) \ .groupBy("city_id", "city_name") \ .agg(F.collect_list("cel_item").alias("cel")) sorted_df.show(2, False)
逻辑拆解
- 拆分数组:同样用
explode把cel数组拆分成单条条目。 - 全局排序:按
city_id、city_name分组排序,同时按subs降序排列,确保同一城市的条目按subs从大到小有序排列。 - 重新聚合:通过
collect_list把同一城市的排序后条目重新聚合成数组,替换原cel列。
两种方式最终都会输出你期望的结果:
+-------+---------+--------------------------------------------+ |city_id|city_name|cel | +-------+---------+--------------------------------------------+ |3273 |city y |[[tlk, 35], [ids, 27], [thr, 24], [smf, 13]]| |3213 |city x |[[thr, 34], [smf, 23], [ids, 17], [tlk, 15]]| +-------+---------+--------------------------------------------+
内容的提问来源于stack exchange,提问作者Wanda
相关产品推荐
相关产品推荐

