在PySpark中修复日期重叠:基于序列与价格的时间范围合并
按接收顺序合并同价格日期范围并处理价格变更的日期截断
原始数据集
+------+------+-------+------------+------------+-----------------+ | key1 | key2 | price | date_start | date_end | sequence_number | +------+------+-------+------------+------------+-----------------+ | a | b | 10 | 2022-01-03 | 2022-01-05 | 1 | | a | b | 10 | 2022-01-02 | 2050-05-15 | 2 | | a | b | 10 | 2022-02-02 | 2022-05-10 | 3 | | a | b | 20 | 2024-02-01 | 2050-10-10 | 4 | | a | b | 20 | 2024-04-01 | 2025-09-10 | 5 | | a | b | 10 | 2022-04-02 | 2024-09-10 | 6 | | a | b | 20 | 2024-09-11 | 2050-10-10 | 7 | +------+------+-------+------------+------------+-----------------+
期望最终结果
+------+------+-------+------------+------------+ | key1 | key2 | price | date_start | date_end | +------+------+-------+------------+------------+ | a | b | 10 | 2022-01-02 | 2024-09-10 | | a | b | 20 | 2024-09-11 | 2050-10-10 | +------+------+-------+------------+------------+
分步处理过程
步骤1:处理前3条序列(sequence_number=1、2、3)
三条记录的key1=a、key2=b、price=10,日期范围存在重叠,合并后取最小date_start和最大date_end:
+------+------+-------+------------+------------+ | key1 | key2 | price | date_start | date_end | +------+------+-------+------------+------------+ | a | b | 10 | 2022-01-02 | 2050-05-15 | +------+------+-------+------------+------------+
步骤2:处理第4条序列(sequence_number=4)
新增price=20的记录,需将原price=10记录的date_end截断为新记录date_start的前一天(2024-01-31):
+------+------+-------+------------+------------+ | key1 | key2 | price | date_start | date_end | +------+------+-------+------------+------------+ | a | b | 10 | 2022-01-02 | 2024-01-31 | | a | b | 20 | 2024-02-01 | 2050-10-10 | +------+------+-------+------------+------------+
步骤3:处理第5条序列(sequence_number=5)
新增的price=20记录日期范围完全被现有price=20的记录覆盖,无需调整:
+------+------+-------+------------+------------+ | key1 | key2 | price | date_start | date_end | +------+------+-------+------------+------------+ | a | b | 10 | 2022-01-02 | 2024-01-31 | | a | b | 20 | 2024-02-01 | 2050-10-10 | +------+------+-------+------------+------------+
步骤4:处理第6条序列(sequence_number=6)
新增price=10的记录,需扩展原price=10记录的date_end至2024-09-10,同时将price=20记录的date_start调整为2024-09-11:
+------+------+-------+------------+------------+ | key1 | key2 | price | date_start | date_end | +------+------+-------+------------+------------+ | a | b | 10 | 2022-01-02 | 2024-09-10 | | a | b | 20 | 2024-09-11 | 2050-10-10 | +------+------+-------+------------+------------+
步骤5:处理第7条序列(sequence_number=7)
新增的price=20记录日期范围与现有price=20的记录完全重合,无需调整:
+------+------+-------+------------+------------+ | key1 | key2 | price | date_start | date_end | +------+------+-------+------------+------------+ | a | b | 10 | 2022-01-02 | 2024-09-10 | | a | b | 20 | 2024-09-11 | 2050-10-10 | +------+------+-------+------------+------------+
内容的提问来源于stack exchange,提问作者CoderWithAGoodName
相关产品推荐
相关产品推荐

