如何合并不同时间间隔的设备时序数据集与错误事件数据集?
问题描述
我们有两个时间序列数据集,需要按照规则合并成包含Time_stamp、Device_IDI、Temperature、Error message的输出数据集:
数据集1:设备10分钟间隔温度记录
| Time_stamp | Device_IDI | Temperature |
|---|---|---|
| 2022-02-18 16:10:00 | 645 | 20 |
| 2022-02-18 16:20:00 | 645 | 21 |
| .... | ||
| 2022-02-18 16:10:00 | 697 | 20 |
| 2022-02-18 16:20:00 | 697 | 20 |
数据集2:设备错误事件记录
| Device_IDI | Start | End | Error message | Duration |
|---|---|---|---|---|
| 645 | 2022-02-18 16:26:05 | 2022-02-19 06:25:17 | temperature error | 13:59:12 |
| 697 | 2022-02-18 16:24:36 | 2022-02-18 16:24:43 | temperature error | 00:00:07 |
合并规则
当设备的10分钟间隔Time_stamp对应的时间段(即Time_stamp到Time_stamp+10分钟)与错误事件的时间区间(Start到End)存在交集时,填充对应的Error message;否则填充Null。
技术实现方案
1. Python Pandas 实现
适合中小规模数据集,步骤清晰易调试:
import pandas as pd # 读取数据集(假设为CSV格式,可根据实际数据源调整) temp_df = pd.read_csv("temperature_data.csv") error_df = pd.read_csv("error_events.csv") # 转换时间列为datetime类型 temp_df['Time_stamp'] = pd.to_datetime(temp_df['Time_stamp']) error_df['Start'] = pd.to_datetime(error_df['Start']) error_df['End'] = pd.to_datetime(error_df['End']) # 计算每个温度记录对应的10分钟时间段 temp_df['interval_start'] = temp_df['Time_stamp'] temp_df['interval_end'] = temp_df['Time_stamp'] + pd.Timedelta(minutes=10) # 按设备ID左连接两个数据集 merged_df = pd.merge(temp_df, error_df, on='Device_IDI', how='left') # 判断时间段交集并填充错误信息 def fill_error(row): if pd.isna(row['Start']): return None # 交集判断逻辑:两个区间[a1,a2]和[b1,b2]存在重叠的条件是 a1 < b2 且 b1 < a2 return row['Error message'] if (row['interval_start'] < row['End']) and (row['Start'] < row['interval_end']) else None merged_df['Error message'] = merged_df.apply(fill_error, axis=1) # 保留目标列并去重(同一记录匹配多个错误事件时取第一个) final_df = merged_df[['Time_stamp', 'Device_IDI', 'Temperature', 'Error message']].drop_duplicates() # 输出或保存结果 print(final_df) # final_df.to_csv("merged_result.csv", index=False)
2. SQL 实现
适合数据库存储的数据集,支持大规模数据查询:
SELECT t.Time_stamp, t.Device_IDI, t.Temperature, CASE WHEN EXISTS ( SELECT 1 FROM error_events e WHERE e.Device_IDI = t.Device_IDI -- 判断温度记录的10分钟区间与错误区间是否重叠 AND t.Time_stamp < e.End AND e.Start < DATE_ADD(t.Time_stamp, INTERVAL 10 MINUTE) ) THEN ( SELECT e.`Error message` FROM error_events e WHERE e.Device_IDI = t.Device_IDI AND t.Time_stamp < e.End AND e.Start < DATE_ADD(t.Time_stamp, INTERVAL 10 MINUTE) LIMIT 1 -- 同一时间段多错误时取第一个,可按需调整 ) ELSE NULL END AS `Error message` FROM temperature_data t ORDER BY t.Device_IDI, t.Time_stamp;
3. Spark SQL 实现
适合超大规模分布式数据集:
SELECT t.Time_stamp, t.Device_IDI, t.Temperature, CASE WHEN EXISTS ( SELECT 1 FROM error_events e WHERE e.Device_IDI = t.Device_IDI AND t.Time_stamp < e.End AND e.Start < date_add(t.Time_stamp, INTERVAL 10 MINUTES) ) THEN ( SELECT e.`Error message` FROM error_events e WHERE e.Device_IDI = t.Device_IDI AND t.Time_stamp < e.End AND e.Start < date_add(t.Time_stamp, INTERVAL 10 MINUTES) LIMIT 1 ) ELSE NULL END AS `Error message` FROM temperature_data t ORDER BY t.Device_IDI, t.Time_stamp;
内容的提问来源于stack exchange,提问作者Thang Trinh
相关产品推荐
相关产品推荐

