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_id | date | close | is_reference | reference_change |
|---|---|---|---|---|
| XXXX | 2000-01-19 | 809.9644 | FALSE | 1.07889737170381 |
| XXXX | 2000-01-20 | 784.274 | FALSE | 1.04467697258748 |
| XXXX | 2000-01-21 | 774.2831 | FALSE | 1.03136878799201 |
| XXXX | 2000-01-24 | 760.0106 | FALSE | 1.0123573811479 |
| XXXX | 2000-01-25 | 750.7335 | FALSE | 1 |
| XXXX | 2000-01-26 | 750.7335 | TRUE | 1 |
| XXXX | 2000-01-27 | 742.17 | FALSE | 0.988593155893536 |
| XXXX | 2000-01-28 | 749.3063 | FALSE | 0.99809892591712 |
| XXXX | 2000-01-31 | 750.02 | FALSE | 0.999049596161621 |
| XXXX | 2000-02-01 | 762.8653 | FALSE | 1.01615992892285 |
| XXXX | 2000-02-02 | 749.3063 | FALSE | 0.99809892591712 |
内容的提问来源于stack exchange,提问作者Eternal Student
相关产品推荐
相关产品推荐

