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

Spark DataFrame基于参考行的窗口值除法计算方案求助

解决方案

针对你的需求,核心是先获取每个symbol_id分组下参考行(is_reference = TRUE)的close值,再将组内每行的close与该参考值做除法计算。如果需要限制仅计算参考行前后n天的数据,也可以额外添加日期范围筛选逻辑。

1. 基础实现:全分组计算reference_change

假设每个symbol_id仅有一条参考记录,先通过窗口函数将参考值广播到组内所有行,再执行除法:

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

// 定义按symbol_id分组的窗口
val symbolPartitionWindow = Window.partitionBy("symbol_id")

// 添加参考值列,提取当前分组中is_reference=true的close值
val dfWithReference = df.withColumn(
  "reference_close",
  first(when(col("is_reference") === true, col("close")), ignoreNulls = true).over(symbolPartitionWindow)
)

// 计算reference_change,删除中间列
val finalDF = dfWithReference.withColumn(
  "reference_change",
  col("close") / col("reference_close")
).drop("reference_close")

2. 进阶:筛选参考行前后n天的数据

如果需要仅保留参考行前后n天的记录(比如n=5),可以先获取参考日期,再通过日期差筛选范围:

// 先提取每个symbol_id的参考日期
val referenceDates = df.filter(col("is_reference") === true)
  .select("symbol_id", "date")
  .withColumnRenamed("date", "reference_date")

// 关联参考日期,计算日期差并筛选前后5天的记录
val dfFilteredByDate = df.join(referenceDates, "symbol_id")
  .withColumn(
    "date_diff",
    datediff(col("date"), col("reference_date"))
  )
  .filter(abs(col("date_diff")) <= 5) // 保留前后5天内的数据

// 计算reference_change,清理中间列
val finalFilteredDF = dfFilteredByDate.withColumn(
  "reference_close",
  first(when(col("is_reference") === true, col("close")), ignoreNulls = true).over(symbolPartitionWindow)
).withColumn(
  "reference_change",
  col("close") / col("reference_close")
).drop("reference_date", "date_diff", "reference_close")

关键逻辑说明

  • first(when(col("is_reference") === true, col("close")), ignoreNulls = true):在每个分组内提取唯一的参考close值,ignoreNulls确保即使分组内参考行位置靠后也能正确获取。
  • 使用datediff计算日期差比rowsBetween更准确,因为实际数据中日期可能存在非连续的情况(比如示例中的周末间隔),行数和天数无法直接对应。

测试输出

将上述代码应用到你的示例数据,会生成你期望的结果:

symbol_iddatecloseis_referencereference_change
XXXX2000-01-19809.9644FALSE1.07889737170381
XXXX2000-01-20784.274FALSE1.04467697258748
XXXX2000-01-21774.2831FALSE1.03136878799201
XXXX2000-01-24760.0106FALSE1.0123573811479
XXXX2000-01-25750.7335FALSE1
XXXX2000-01-26750.7335TRUE1
XXXX2000-01-27742.17FALSE0.988593155893536
XXXX2000-01-28749.3063FALSE0.99809892591712
XXXX2000-01-31750.02FALSE0.999049596161621
XXXX2000-02-01762.8653FALSE1.01615992892285
XXXX2000-02-02749.3063FALSE0.99809892591712

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 12:30:43