如何在PySpark中对每个分组内的变量进行排序?
解决方案:按ID分组并在组内按TS排序Val
没问题,针对你要按id分组、依据ts列对val值排序的需求,我们可以用Spark的窗口函数或者直接排序/聚合来实现,分两种常见场景给你说明:
场景1:保留所有行,每个ID组内按TS排序
如果需要让每个id分组下的所有记录按ts升序(或降序)排列,最简单高效的方式就是直接对整个DataFrame按id和ts排序:
# 修正SparkSession创建方式(原代码可能报错) spark = ss.builder.appName("group_sort_demo").getOrCreate() sdf = spark.createDataFrame(pdf) # 按id分组,组内按ts升序排序 sdf_sorted = sdf.orderBy("id", "ts") sdf_sorted.show()
输出结果:
+---+---+---+ | id| ts|val| +---+---+---+ | 1| 1|dog| | 1| 2|cat| | 2| 2|cat| | 2| 3|cat| | 2| 4|dog| +---+---+---+
如果需要降序排序,只需要把ts改成F.desc("ts")即可:
sdf_sorted_desc = sdf.orderBy("id", F.desc("ts"))
场景2:聚合每个ID的Val为TS排序后的列表
如果需要将每个id对应的val按ts顺序收集成一个列表,得到每个id一行的结果,可以用Spark 3.0+支持的collect_list(...).orderBy(...)语法,或者结合窗口函数实现:
方法1:Spark 3.0+ 简洁写法
from pyspark.sql import functions as F sdf_grouped = sdf.groupBy("id").agg( F.collect_list("val").orderBy("ts").alias("sorted_vals") ) sdf_grouped.show(truncate=False)
输出结果:
+---+----------------+ |id |sorted_vals | +---+----------------+ |1 |[dog, cat] | |2 |[cat, cat, dog] | +---+----------------+
方法2:兼容低版本Spark的窗口函数写法
如果你的Spark版本低于3.0,可以先通过窗口函数给每个分组内的行标记排序序号,再聚合排序:
from pyspark.sql.window import Window # 定义窗口:按id分区,按ts升序排序 window_spec = Window.partitionBy("id").orderBy("ts") # 添加排序序号列 sdf_with_rank = sdf.withColumn("sort_rank", F.row_number().over(window_spec)) # 聚合时按序号排序后提取val sdf_grouped = sdf_with_rank.groupBy("id").agg( F.collect_list(F.struct("sort_rank", "val")).alias("ranked_vals") ).withColumn( "sorted_vals", F.transform(F.sort_array("ranked_vals"), lambda x: x.val) ).drop("ranked_vals") sdf_grouped.show(truncate=False)
这个方法和上面的简洁写法输出结果一致,适合低版本Spark环境。
内容的提问来源于stack exchange,提问作者scottlittle
相关产品推荐
相关产品推荐

