在Databricks中将重复OS版本数据转换为时间区间的方案
实现设备OS版本生效时间区间转换(Databricks Delta)
问题背景
我们需要处理每日上报的设备OS版本数据(无论版本是否变更),将其转换为每个设备对应OS版本的生效时间区间:
- 同一设备连续上报的相同版本需合并为一个时间区间
- 版本变更(包括降级)时生成新的区间
- 最新版本的
EndTime固定为2999-12-31 23:59:59
假设原始Delta表device_os_daily结构如下:
| device_id | report_time | os_version |
|---|---|---|
| Device_A | 2024-01-01 00:00:00 | OS_1 |
| Device_A | 2024-01-02 00:00:00 | OS_1 |
| Device_A | 2024-01-03 00:00:00 | OS_2 |
| Device_B | 2024-01-01 00:00:00 | OS_3 |
| Device_B | 2024-01-02 00:00:00 | OS_1 |
PySpark 实现方案
from pyspark.sql import Window from pyspark.sql.functions import col, lag, when, min, lead, lit # 读取原始Delta表 raw_df = spark.read.format("delta").table("device_os_daily") # 1. 按设备分组、上报时间排序,标记连续相同版本的分组 device_window = Window.partitionBy("device_id").orderBy("report_time") grouped_df = raw_df.withColumn( "is_new_version", when(lag("os_version").over(device_window) != col("os_version"), 1) .when(lag("os_version").over(device_window).isNull(), 1) .otherwise(0) ).withColumn( "version_group", sum("is_new_version").over(device_window.rowsBetween(Window.unboundedPreceding, Window.currentRow)) ) # 2. 计算每个版本组的起始时间 version_group_window = Window.partitionBy("device_id", "version_group") interval_df = grouped_df.withColumn( "start_time", min("report_time").over(version_group_window) ).select("device_id", "os_version", "version_group", "start_time").distinct() # 3. 获取下一个版本的起始时间,计算当前版本的结束时间 next_version_window = Window.partitionBy("device_id").orderBy("start_time") result_df = interval_df.withColumn( "next_start_time", lead("start_time").over(next_version_window) ).withColumn( "end_time", when(col("next_start_time").isNull(), lit("2999-12-31 23:59:59").cast("timestamp")) .otherwise(col("next_start_time") - lit(1).cast("interval day")) ).select("device_id", "os_version", "start_time", "end_time") # 写入结果Delta表 result_df.write.format("delta").mode("overwrite").saveAsTable("device_os_version_intervals")
代码说明
- 通过
lag函数判断当前版本与上一次上报版本是否一致,生成分组标记 - 用累加分组标记的方式,将连续相同的版本归为同一组
- 取每组最早的上报时间作为
start_time,用lead获取下一组的起始时间,计算当前组的end_time(下一组起始时间减1天) - 无后续版本的组,
end_time设为固定值
Databricks SQL 实现方案
CREATE OR REPLACE TABLE device_os_version_intervals USING DELTA AS WITH version_groups AS ( -- 标记连续相同版本的分组 SELECT device_id, os_version, report_time, SUM(CASE WHEN LAG(os_version) OVER (PARTITION BY device_id ORDER BY report_time) != os_version OR LAG(os_version) OVER (PARTITION BY device_id ORDER BY report_time) IS NULL THEN 1 ELSE 0 END) OVER (PARTITION BY device_id ORDER BY report_time ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW) AS version_group FROM device_os_daily ), version_starts AS ( -- 计算每个版本组的起始时间 SELECT device_id, os_version, version_group, MIN(report_time) AS start_time FROM version_groups GROUP BY device_id, os_version, version_group ), interval_calculation AS ( -- 获取下一个版本的起始时间 SELECT device_id, os_version, start_time, LEAD(start_time) OVER (PARTITION BY device_id ORDER BY start_time) AS next_start_time FROM version_starts ) SELECT device_id, os_version, start_time, CASE WHEN next_start_time IS NULL THEN TIMESTAMP('2999-12-31 23:59:59') ELSE next_start_time - INTERVAL 1 DAY END AS end_time FROM interval_calculation ORDER BY device_id, start_time;
性能优化建议
- 对原始Delta表按
device_id或report_time分区,减少数据扫描范围 - 若上报时间为任意时间戳,将
END_TIME的计算逻辑改为next_start_time - INTERVAL 1 MICROSECOND,确保时间区间无重叠 - 对于超大规模数据,可添加
OPTIMIZE和ZORDER BY命令优化Delta表:OPTIMIZE device_os_version_intervals ZORDER BY device_id;
内容的提问来源于stack exchange,提问作者khidir sanosi
相关产品推荐
相关产品推荐

