PySpark按ID分组按时间戳排序聚合列成有序列表问题排查
代码问题排查
原代码无法得到预期结果,核心问题有3个:
- 字段引用错误:代码中
collect_list传入的字段名是reco、score,但原始DataFrame中不存在这两个字段,对应字段实际为col1、col2,同时漏了timestamp字段的聚合 - 窗口范围配置错误:代码中定义了全分区范围的
ranged_spec但未实际使用,带orderBy的窗口默认范围是分区起点到当前行,会生成逐行递增的累计列表,而非整个分组的全量列表 - 输出逻辑冗余:窗口函数会为分组内每一行都生成相同的聚合结果,最终需要去重才能得到每个Id一行的输出,性能远低于分组聚合方案
最优实现方案(原生函数,性能最高)
使用groupBy+结构体排序+列表收集的方式实现,全程用Spark内置函数,无UDF性能损耗,能严格保证列表按timestamp升序排列:
from pyspark.sql import functions as F df_result = df.groupBy("Id") \ # 打包同组下的三个字段为结构体,收集后按timestamp升序排序 .agg(F.sort_array(F.collect_list(F.struct("timestamp", "col1", "col2"))).alias("sorted_data")) \ # 从排序后的结构体中拆分出三个字段的有序列表 .select( "Id", F.col("sorted_data.timestamp").alias("timestamp"), F.col("sorted_data.col1").alias("col1"), F.col("sorted_data.col2").alias("col2") )
执行后输出完全匹配预期:
+---+----------+------+------+ | Id| timestamp| col1| col2| +---+----------+------+------+ |abc|[123, 789]|[1, 0]|[0, 1]| |def|[321, 456]|[0, 1]|[1, 0]| +---+----------+------+------+
窗口函数修正写法(不推荐,仅作参考)
如果一定要用窗口函数实现,需要修正窗口范围、字段名,最后加去重逻辑:
from pyspark.sql import functions as F from pyspark.sql import Window as W # 定义分区+排序规则 window_spec = W.partitionBy("Id").orderBy("timestamp") # 明确指定窗口范围为整个分区 full_partition_spec = window_spec.rowsBetween(W.unboundedPreceding, W.unboundedFollowing) df_result = df.withColumn("timestamp", F.collect_list("timestamp").over(full_partition_spec))\ .withColumn("col1", F.collect_list("col1").over(full_partition_spec))\ .withColumn("col2", F.collect_list("col2").over(full_partition_spec))\ .select("Id", "timestamp", "col1", "col2")\ .dropDuplicates(["Id"]) # 去除组内重复的聚合结果,每个Id仅留一行
注意:该方案会为分组内所有行重复计算全量聚合列表,数据量较大时性能明显差于分组聚合方案,非必要不使用。
内容的提问来源于stack exchange,提问作者Abhishek Patil
相关产品推荐
相关产品推荐

