基于df_purchase日期排名为df_obs的Date_Observed列分配对应排名
Spark DataFrame:为观测日期匹配对应采购日期区间的Rank值
现有数据
1. df_obs(观测数据)
执行df_obs.show()后输出:
| Item No | Date_Observed |
|---|---|
| Item 1 | 2021-09-20 |
| Item 1 | 2022-12-05 |
| Item 2 | 2022-10-27 |
| Item 1 | 2022-09-20 |
| Item 2 | 2023-02-20 |
| Item 2 | 2023-03-20 |
| Item 1 | 2023-01-20 |
2. df_purchase(已生成Rank的采购数据)
执行以下代码生成rank列:
df_purchase.withColumn("rank", dense_rank().over(Window.partitionBy("ITEM_No").orderBy(asc("Date_Purchase")))).show()
输出结果:
| Item_No | Date_Purchase | rank |
|---|---|---|
| Item 1 | 2021-08-21 | 1 |
| Item 1 | 2022-02-23 | 2 |
| Item 1 | 2022-12-29 | 3 |
| Item 2 | 2022-09-20 | 1 |
| Item 2 | 2023-01-20 | 2 |
需求说明
根据每个Item对应的采购日期区间,为df_obs的每条观测记录分配对应的rank:
- 若观测日期落在某两次采购日期之间(即
Date_Purchase[rank=n] ≤ Date_Observed < Date_Purchase[rank=n+1]),则匹配rank=n - 若观测日期晚于该Item的最后一次采购日期,则匹配该Item的最大rank值
示例:df_obs中Item 1的2022-12-05观测日期,落在rank=2(2022-02-23)和rank=3(2022-12-29)的区间内,因此分配rank=2。
期望输出
| Item No | Date_Observed | rank |
|---|---|---|
| Item 1 | 2021-09-20 | 1 |
| Item 1 | 2022-12-05 | 2 |
| Item 2 | 2022-10-27 | 1 |
| Item 1 | 2022-09-20 | 2 |
| Item 2 | 2023-02-20 | 2 |
| Item 2 | 2023-03-20 | 2 |
| Item 1 | 2023-01-20 | 3 |
解决方案
步骤说明
- 为
df_purchase添加next_purchase_date列,标记该Item下一次采购的日期(使用lead窗口函数) - 将
df_obs与处理后的df_purchase按Item关联,筛选观测日期落在当前采购日期和下一次采购日期之间的记录;对于没有下一次采购日期的最大rank记录,直接匹配所有晚于当前采购日期的观测记录 - 整理输出列,得到最终结果
代码实现(Scala)
import org.apache.spark.sql.functions._ import org.apache.spark.sql.Window // 1. 处理df_purchase,添加下一次采购日期列 val purchaseWithNextDate = df_purchase.withColumn( "next_purchase_date", lead("Date_Purchase", 1).over(Window.partitionBy("Item_No").orderBy("Date_Purchase")) ) // 2. 关联df_obs和处理后的采购数据,匹配对应的rank val result = df_obs.join( purchaseWithNextDate, df_obs("Item No") === purchaseWithNextDate("Item_No") && (df_obs("Date_Observed") >= purchaseWithNextDate("Date_Purchase") && (purchaseWithNextDate("next_purchase_date").isNull || df_obs("Date_Observed") < purchaseWithNextDate("next_purchase_date"))), "left" ).select( df_obs("Item No"), df_obs("Date_Observed"), purchaseWithNextDate("rank") ) // 查看结果 result.show()
补充说明
如果你的日期列是字符串类型,需要先转换为Date类型:
df_obs = df_obs.withColumn("Date_Observed", to_date(col("Date_Observed"), "yyyy-MM-dd")) df_purchase = df_purchase.withColumn("Date_Purchase", to_date(col("Date_Purchase"), "yyyy-MM-dd"))
内容的提问来源于stack exchange,提问作者imtig
相关产品推荐
相关产品推荐

