Databricks Delta-Tempo中interpol与resample实现问题:时序聚合异常
处理Databricks Delta-Tempo时序聚合问题
问题背景
使用Databricks Delta-Tempo处理秒级时序数据,目标是将数据聚合至每日小时级别及每日级别。数据集时间字段为字符串类型,Pressure为double类型,示例数据如下:
| Time (String) | Pressure |
|---|---|
| 02-01-2023 00:00:00 | 2720.78 |
| 02-01-2023 00:00:01 | 2720.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')
结果不符合预期。
解决方案
问题根源分析
- 分区列设计冲突:
PartitionDayHourMin用日+小时拼接,Delta-Tempo的分区列用于分组时序处理,该列将同一日+小时的数据归为一组,与按day/hr的聚合频率冲突,导致时间戳生成异常。 - 方法误用:
interpolate核心作用是缺失值插值,不是时序聚合,降采样聚合应使用resample或PySpark原生窗口函数。 - 时间过滤不匹配:原始时间转换后为
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
相关产品推荐
相关产品推荐

