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

Databricks Delta-Tempo中interpol与resample实现问题:时序聚合异常

处理Databricks Delta-Tempo时序聚合问题

问题背景

使用Databricks Delta-Tempo处理秒级时序数据,目标是将数据聚合至每日小时级别及每日级别。数据集时间字段为字符串类型,Pressure为double类型,示例数据如下:

Time (String)Pressure
02-01-2023 00:00:002720.78
02-01-2023 00:00:012720.7

已尝试操作及问题

1. 时间转换与分区列创建

通过PySpark将时间列转换为timestamp类型,并创建PartitionDayHourMin分区列,代码如下:

csvdata_df = spark.read.csv("resources\PTGauge.csv", inferSchema=True, header=True)\
    .withColumn("TimeStampUTC", to_timestamp("Time (Utc)","MM-dd-yyyy HH:mm:ss"))\
    .withColumn("PartitionDayHourMin", concat(lpad(dayofmonth("TimeStampUTC"),2,"0"), lpad(hour("TimeStampUTC"),2,"0").cast("string")))

转换后的DataFrame示例:

+-------------------+--------+-------------------+-------------------+
|      Time_Utc      |Pressure|     TimeStampUTC   | PartitionDayHourMin|
+-------------------+--------+-------------------+-------------------+
|02-01-2023 00:00:00| 2720.78|2023-01-02 00:00:00|               0200|
|02-01-2023 00:00:01|  2720.7|2023-01-02 00:00:01|               0200|
+-------------------+--------+-------------------+-------------------+

2. interpolate方法的问题

创建TSDF后调用interpolate方法:

tempo_tsdf = TSDF(csvdata_df, ts_col="TimeStampUTC", partition_cols = ["PartitionDayHourMin"])
interpolated_tsdf1 = tempo_tsdf.interpolate(
    freq="day",
    func="mean",
    target_cols= ["Pressure"],
    method="linear"
)
print(">>> Using TDSF interpolate show")
interpolated_tsdf1.df.show(100)

遇到的问题:

  • 结果中PartitionDayHourMin按日和小时划分,但TimeStampUTC未显示对应时段的正确值,无法用于绘图
  • 使用freq="hr"时无法得到预期的小时级数据,仅freq="day"有结果

3. resample方法的问题

尝试用resample方法:

start_ts = '2023-02-01T00:00:00.000+0000'
end_ts = '2023-02-01T23:59:59.000+0000'
interval_inclusive = tempo_tsdf.between(start_ts, end_ts)
resampled_sdf = interval_inclusive.resample(freq='hr', func='mean')

结果不符合预期。

解决方案

问题根源分析

  1. 分区列设计冲突:PartitionDayHourMin用日+小时拼接,Delta-Tempo的分区列用于分组时序处理,该列将同一日+小时的数据归为一组,与按day/hr的聚合频率冲突,导致时间戳生成异常。
  2. 方法误用:interpolate核心作用是缺失值插值,不是时序聚合,降采样聚合应使用resample或PySpark原生窗口函数。
  3. 时间过滤不匹配:原始时间转换后为2023-01-02,但过滤范围设为2023-02-01,无匹配数据导致结果异常。

修正步骤

1. 调整分区逻辑

若数据量不大,可直接去掉不合理的分区列:

tempo_tsdf = TSDF(csvdata_df, ts_col="TimeStampUTC")

若需分区,改用日期字符串作为分区列(如yyyy-MM-dd),避免与聚合频率冲突:

from pyspark.sql.functions import date_format

csvdata_df = csvdata_df.withColumn("PartitionDate", date_format("TimeStampUTC", "yyyy-MM-dd"))
tempo_tsdf = TSDF(csvdata_df, ts_col="TimeStampUTC", partition_cols=["PartitionDate"])

2. 正确使用resample做降采样聚合

Delta-Tempo的resample需使用标准freq参数值,且确保时间范围匹配数据实际区间:

# 匹配数据实际起始时间
start_ts = '2023-01-02T00:00:00.000+0000'
end_ts = '2023-01-02T23:59:59.000+0000'

filtered_tsdf = tempo_tsdf.between(start_ts, end_ts)
# 小时级聚合(freq用"hour"而非"hr")
hourly_resampled = filtered_tsdf.resample(freq='hour', func='mean', target_cols=["Pressure"])
# 每日级聚合
daily_resampled = filtered_tsdf.resample(freq='day', func='mean', target_cols=["Pressure"])

# 查看结果
hourly_resampled.df.orderBy("TimeStampUTC").show()
daily_resampled.df.orderBy("TimeStampUTC").show()

3. 替代方案:PySpark原生窗口函数聚合

若Delta-Tempo方法仍有问题,可直接用PySpark原生分组聚合,时间戳更直观:

from pyspark.sql.functions import date_trunc, avg

# 小时级聚合
hourly_df = csvdata_df.groupBy(date_trunc("hour", "TimeStampUTC").alias("HourStart"))\
    .agg(avg("Pressure").alias("AvgPressure"))

# 每日级聚合
daily_df = csvdata_df.groupBy(date_trunc("day", "TimeStampUTC").alias("DayStart"))\
    .agg(avg("Pressure").alias("AvgPressure"))

# 按时间排序查看结果
hourly_df.orderBy("HourStart").show()
daily_df.orderBy("DayStart").show()

关键注意点

  • 确认时间转换正确性:原始时间格式为MM-dd-yyyy,转换后需验证TimeStampUTC与实际日期对应,避免时区或格式解析错误。
  • 区分方法用途:interpolate用于补全缺失时序点,聚合需用resample或原生分组函数。
  • 分区列与聚合频率需兼容,避免分组逻辑冲突。

内容的提问来源于stack exchange,提问作者rbpjava

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 00:32:10