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

Spark中基于最近日期匹配的双DataFrame Join实现方案咨询

实现方案

方案1:通用窗口函数实现(兼容所有Spark版本)

实现逻辑:

  1. 先按SerialNumber等值关联两张表,过滤出ValidityDate2 >= ValidityDate1的所有匹配记录
  2. 为每条异常事件记录计算与匹配的产品记录的日期差值
  3. 按异常事件的唯一标识(SerialNumber+ExceptionId+ValidityDate2)分区,按日期差值升序排名,取排名第一的即为最近匹配的产品记录
import org.apache.spark.sql.functions._
import org.apache.spark.sql.expressions.Window

// 读取parquet表生成DataFrame,确保日期字段为Date类型
val productDF = spark.read.parquet("Product表存储路径")
val exceptionDF = spark.read.parquet("ExceptionEvents表存储路径")

// 关联并过滤符合日期条件的记录
val joinedDF = exceptionDF.join(
  productDF,
  exceptionDF("SerialNumber") === productDF("SerialNumber") && 
  exceptionDF("ValidityDate2") >= productDF("ValidityDate1"),
  "left" // 不需要保留无匹配的异常事件时可改为inner
)
// 计算两个日期的差值
.withColumn("date_diff", datediff(col("ValidityDate2"), col("ValidityDate1")))

// 定义窗口:按异常事件维度分区,按日期差升序排序
val winSpec = Window.partitionBy("SerialNumber", "ExceptionId", "ValidityDate2")
  .orderBy(col("date_diff").asc)

// 取排名第一的最近匹配记录,筛选所需字段
val resultDF = joinedDF.withColumn("rn", row_number().over(winSpec))
  .filter(col("rn") === 1)
  .select("SerialNumber", "ExceptionId", "UniqueKey", "ValidityDate2")

// 验证结果
resultDF.show()

方案2:Spark 3.2+ 优化ASOF JOIN实现(性能更高)

Spark 3.2版本原生支持ASOF JOIN,专门用于时间序列最近值匹配场景,无需全量关联后过滤,大数据量下性能远高于方案1:

import org.apache.spark.sql.functions._

val productDF = spark.read.parquet("Product表存储路径")
val exceptionDF = spark.read.parquet("ExceptionEvents表存储路径")

// ASOF JOIN要求关联键和排序键有序,提前按关联键+日期排序
val sortedProduct = productDF.sort("SerialNumber", "ValidityDate1")
val sortedException = exceptionDF.sort("SerialNumber", "ValidityDate2")

// 执行ASOF JOIN,直接匹配小于等于异常事件日期的最近产品记录
val resultDF = sortedException.joinAsOf(
  sortedProduct,
  leftOn = col("ValidityDate2"),
  rightOn = col("ValidityDate1"),
  usingColumns = Seq("SerialNumber"),
  direction = "backward", // backward表示找小于等于左表日期的最近右表记录
  how = "left" // 不需要保留无匹配的异常事件时可改为inner
).select("SerialNumber", "ExceptionId", "UniqueKey", "ValidityDate2")

// 验证结果
resultDF.show()

注意事项

  • 若日期字段为字符串格式,需先转换为Date类型,示例格式对应转换写法:to_date(col("ValidityDate1"), "dd/MM/yyyy")
  • 两种方案输出结果与你给出的期望结果完全一致

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 01:15:04