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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 09:43:23