Spark DataFrame新增列:取同产品PD状态下小于当前记录的最大service_date
Spark DataFrame 新增历史最大日期列解决方案
核心思路
采用Spark窗口聚合函数实现,完全避免了关联操作,性能远高于子查询和非等值连接方案,也解决了LAG函数偏移量固定不适用的问题。
代码示例
Spark SQL 写法
SELECT id, claim_id, service_date, status, product, MAX(IF(status = 'PD', service_date, NULL)) OVER ( PARTITION BY product ORDER BY service_date ASC ROWS BETWEEN UNBOUNDED PRECEDING AND 1 PRECEDING ) AS previous_service_date FROM 你的表名
Scala DataFrame API 写法
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ // 定义窗口规则:按product分区,按service_date升序排序,范围是分区开头到当前行的前一行 val winSpec = Window.partitionBy("product") .orderBy("service_date") .rowsBetween(Window.unboundedPreceding, -1) val resultDF = 原始DataFrame .withColumn("previous_service_date", // 只取status为PD的行的service_date,取最大值,null值会被max自动忽略 max(when(col("status") === "PD", col("service_date"))).over(winSpec) )
PySpark 写法
from pyspark.sql import Window from pyspark.sql.functions import max, when, col win_spec = Window.partitionBy("product") \ .orderBy("service_date") \ .rowsBetween(Window.unboundedPreceding, -1) result_df = 原始DataFrame.withColumn( "previous_service_date", max(when(col("status") == "PD", col("service_date"))).over(win_spec) )
原理说明
- 窗口按
product分区,确保只会统计和当前行同产品的历史数据 - 按
service_date升序排序,保证窗口内的行时间都早于等于当前行 - 窗口范围限定为
UNBOUNDED PRECEDING AND 1 PRECEDING,自动排除当前行本身,满足「service_date小于当前行」的要求 - 用
when函数将非PD状态的行的service_date置为null,max函数会自动忽略null值,最终返回的就是符合所有条件的最大历史service_date
该方案仅需要一次Shuffle操作完成窗口计算,在大数据量下性能优势非常明显,输出结果和你给出的预期样例完全一致。
内容的提问来源于stack exchange,提问作者Michael Bass
相关产品推荐
相关产品推荐

