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

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)
)

原理说明

  1. 窗口按product分区,确保只会统计和当前行同产品的历史数据
  2. 按service_date升序排序,保证窗口内的行时间都早于等于当前行
  3. 窗口范围限定为UNBOUNDED PRECEDING AND 1 PRECEDING,自动排除当前行本身,满足「service_date小于当前行」的要求
  4. 用when函数将非PD状态的行的service_date置为null,max函数会自动忽略null值,最终返回的就是符合所有条件的最大历史service_date

该方案仅需要一次Shuffle操作完成窗口计算,在大数据量下性能优势非常明显,输出结果和你给出的预期样例完全一致。

内容的提问来源于stack exchange,提问作者Michael Bass

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 12:15:05