You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在Databricks中将重复OS版本数据转换为时间区间的方案

实现设备OS版本生效时间区间转换(Databricks Delta)

问题背景

我们需要处理每日上报的设备OS版本数据(无论版本是否变更),将其转换为每个设备对应OS版本的生效时间区间:

  • 同一设备连续上报的相同版本需合并为一个时间区间
  • 版本变更(包括降级)时生成新的区间
  • 最新版本的EndTime固定为2999-12-31 23:59:59

假设原始Delta表device_os_daily结构如下:

device_idreport_timeos_version
Device_A2024-01-01 00:00:00OS_1
Device_A2024-01-02 00:00:00OS_1
Device_A2024-01-03 00:00:00OS_2
Device_B2024-01-01 00:00:00OS_3
Device_B2024-01-02 00:00:00OS_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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.12 00:37:31