如何调整PySpark窗口函数实现按日滞后统计预约到场/缺席数
解决PySpark DataFrame按日期滞后统计到场/缺席数的问题
你的问题核心是原窗口函数会包含同日期的前置行,导致统计范围不符合要求。要实现仅统计当前行日期之前的记录(同日期行不互相统计),可以用以下两种方案:
方案一:先按日聚合再关联累计值(直观易维护)
这种方法先按用户和日期汇总每日的到场/缺席数,再计算截止到前一天的累计值,最后关联回原表,确保同日期行共享相同的前置统计结果。
代码步骤:
from pyspark.sql import functions as F from pyspark.sql import Window # 创建样本DataFrame并转换日期类型 data = [ (1, '2023-01-01', True), (1, '2023-01-01', True), (1, '2023-01-02', True), (1, '2023-01-02', False), (1, '2023-01-03', False), (1, '2023-01-04', True), (2, '2023-01-02', True), (2, '2023-01-02', False), (2, '2023-01-03', False), ] columns = ['psyin_iden_rn', 'psdln.dat', 'psdln.awz'] df = spark.createDataFrame(data, columns) df = df.withColumn('psdln.dat', F.to_date(F.col('psdln.dat'))) # 1. 按用户+日期聚合,计算每日到场/缺席数 daily_agg = df.groupBy('psyin_iden_rn', 'psdln.dat').agg( F.sum(F.when(F.col('psdln.awz') == True, 1).otherwise(0)).alias('daily_show'), F.sum(F.when(F.col('psdln.awz') == False, 1).otherwise(0)).alias('daily_no_show') ) # 2. 计算每个用户截止到前一天的累计到场/缺席数 window_daily = Window.partitionBy('psyin_iden_rn').orderBy('psdln.dat').rowsBetween(Window.unboundedPreceding, -1) daily_agg_with_prev = daily_agg.withColumn( 'prev_show_cnt', F.coalesce(F.sum('daily_show').over(window_daily), F.lit(0)) ).withColumn( 'prev_no_show_cnt', F.coalesce(F.sum('daily_no_show').over(window_daily), F.lit(0)) ) # 3. 关联回原表,得到每条记录的前置统计值 result_df = df.join( daily_agg_with_prev.select('psyin_iden_rn', 'psdln.dat', 'prev_show_cnt', 'prev_no_show_cnt'), on=['psyin_iden_rn', 'psdln.dat'], how='left' ).withColumnRenamed('prev_show_cnt', 'show_cnt').withColumnRenamed('prev_no_show_cnt', 'no_show_cnt') # 查看结果 result_df.show()
方案二:直接使用rangeBetween窗口函数(更简洁)
利用日期转时间戳的特性,通过rangeBetween限定统计范围为当前日期前一天及更早,避免包含同日期的行。
代码步骤:
from pyspark.sql import functions as F from pyspark.sql import Window # 创建样本DataFrame并转换日期类型 data = [ (1, '2023-01-01', True), (1, '2023-01-01', True), (1, '2023-01-02', True), (1, '2023-01-02', False), (1, '2023-01-03', False), (1, '2023-01-04', True), (2, '2023-01-02', True), (2, '2023-01-02', False), (2, '2023-01-03', False), ] columns = ['psyin_iden_rn', 'psdln.dat', 'psdln.awz'] df = spark.createDataFrame(data, columns) df = df.withColumn('psdln.dat', F.to_date(F.col('psdln.dat'))) # 定义窗口:按用户分组,按日期时间戳排序,范围限定为当前日期前一天及更早 window_spec = Window.partitionBy('psyin_iden_rn').orderBy(F.unix_timestamp('psdln.dat')).rangeBetween( Window.unboundedPreceding, F.unix_timestamp('psdln.dat') - 86400 # 86400秒=1天 ) # 计算到场/缺席数 result_df = df.withColumn( 'show_cnt', F.coalesce(F.sum(F.when(F.col('psdln.awz') == True, 1).otherwise(0)).over(window_spec), F.lit(0)) ).withColumn( 'no_show_cnt', F.coalesce(F.sum(F.when(F.col('psdln.awz') == False, 1).otherwise(0)).over(window_spec), F.lit(0)) ) # 查看结果 result_df.show()
结果验证
两种方案都会得到符合预期的输出:
- 例如用户1的2023-01-01的两行记录,
show_cnt和no_show_cnt都是0(没有更早的记录) - 用户1的2023-01-02的两行记录,
show_cnt=2(来自2023-01-01的两个True),no_show_cnt=0 - 用户2的2023-01-03的记录,
show_cnt=1,no_show_cnt=1(来自2023-01-02的统计)
内容的提问来源于stack exchange,提问作者Mitchell
相关产品推荐
相关产品推荐

