PySpark DataFrame带特殊条件的跨行日期列计算求助
解决PySpark DataFrame按条件跨行计算日期列的问题
没问题,咱们可以借助PySpark的窗口函数和条件判断逻辑来实现这个需求,正好你已经确认数据是正确排序的,这是实现的关键前提。
完整实现代码
from pyspark.sql import Window from pyspark.sql.functions import expr, lag, when # 定义窗口规范:基于已排序的数据,获取前一行的date值 window_spec = Window.orderBy("date") # 计算目标new_date列,先临时生成前一行日期的辅助列 result_df = dataframe.withColumn( "prev_date", lag("date").over(window_spec) ).withColumn( "new_date", expr("date_add(CASE WHEN flag = 1 THEN date ELSE prev_date END, diff)") ).drop("prev_date") # 清理临时辅助列
代码细节解释
- 窗口函数
lag:用来提取当前行的前一行date值,存入临时列prev_date。因为你已经保证数据排序正确,Window.orderBy("date")能确保我们取到的是逻辑上的前一行数据。 - 条件判断
CASE WHEN:在date_add计算里,根据flag的值动态选择日期基准:- 当
flag=1时,直接用当前行的date与diff相加 - 其他情况,用上一行的
prev_date与当前行的diff相加
- 当
- 清理临时列:
prev_date只是计算过程中的辅助列,最后用drop删掉即可,不影响最终结果。
示例数据验证结果
用你提供的测试数据运行后,会得到符合预期的new_date列:
| flag | date | diff | new_date |
|---|---|---|---|
| 1 | 2014-05-31 | 0 | 2014-05-31 |
| 2 | 2014-06-02 | 2 | 2014-06-02 |
| 3 | 2016-01-14 | 591 | 2017-08-27 |
| 1 | 2016-07-08 | 0 | 2016-07-08 |
| 2 | 2016-07-12 | 4 | 2016-07-12 |
内容的提问来源于stack exchange,提问作者MaryJane
相关产品推荐
相关产品推荐

