You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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)

逻辑拆解

  1. 拆分数组:F.explode("cel")把每个城市对应的cel数组拆分成单独行,每行对应一个(carr, subs)结构体,命名为cel_item。
  2. 窗口排序收集:窗口的partitionBy确保我们只在同一个城市的范围内处理数据,orderBy指定按subs降序排列;再通过collect_list在窗口范围内收集排序后的结构体,生成每个城市的排序后数组。
  3. 恢复原结构:因为同一个城市的所有拆分行会生成完全相同的排序后数组,用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)

逻辑拆解

  1. 拆分数组:同样用explode把cel数组拆分成单条条目。
  2. 全局排序:按city_id、city_name分组排序,同时按subs降序排列,确保同一城市的条目按subs从大到小有序排列。
  3. 重新聚合:通过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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.30 05:47:32