PySpark对数据子集应用窗口函数且保留全量数据的实现
问题根因
原有代码丢失Q=1、Q=4行的核心原因是在窗口计算前调用了where做过滤,不符合条件的行在计算环节就已经被从DataFrame中移除,后续计算自然无法保留全量原始数据。
正确实现逻辑是不提前过滤行,在计算平均值时仅提取Q=2、Q=3行的median值参与运算,avg函数会自动忽略null值,既保证计算逻辑正确,也能保留所有原始行。
具体实现
前置依赖与测试数据
from pyspark.sql import functions as F from pyspark.sql.window import Window # 构造原始测试数据 df = spark.createDataFrame([ ('2018-03-31',6,1),('2018-03-31',27,2),('2018-03-31',3,3),('2018-03-31',44,4), ('2018-06-30',6,1),('2018-06-30',4,3),('2018-06-30',32,2),('2018-06-30',112,4), ('2018-09-30',2,1),('2018-09-30',23,4),('2018-09-30',37,3),('2018-09-30',3,2) ],['date','median','Q']) # 定义按date分区的窗口,分区内全量聚合不需要指定orderBy window = Window.partitionBy("date")
方案1:同日期所有行均展示result(匹配第二种预期输出)
直接在窗口计算时用条件判断筛选参与运算的值,计算结果会覆盖分区内所有行:
df_all_show = df.withColumn( "result", F.avg( F.when(F.col("Q").isin(2, 3), F.col("median")) ).over(window) ) df_all_show.show()
逻辑说明:
F.when会将Q不等于2、3的行的计算入参转为null,avg聚合时自动跳过null值,最终得到的就是每个date下仅Q=2、Q=3行的median平均值- 没有提前过滤行,所有原始数据都会保留,同一date下所有行的result值完全一致
方案2:仅Q=2、Q=3行展示result,其余行result为null(匹配第一种预期输出)
在方案1的计算逻辑基础上,再加一层条件判断,非目标Q值行的result字段置为null:
df_target_show = df.withColumn( "result", F.when( F.col("Q").isin(2, 3), F.avg(F.when(F.col("Q").isin(2, 3), F.col("median"))).over(window) ) ) df_target_show.show()
结果校验
手动计算各日期的平均值,和代码输出完全匹配:
- 2018-03-31:(27+3)/2 = 15
- 2018-06-30:(32+4)/2 = 18
- 2018-09-30:(3+37)/2 = 20
内容的提问来源于stack exchange,提问作者cnns
相关产品推荐
相关产品推荐

