基于时间序列与日志数据生成移动求和列的实现方案问询
问题:时间序列与非连续日志数据合并及2天移动求和实现
数据生成代码
时间序列数据
import pandas as pd import numpy as np df = pd.DataFrame(index=pd.date_range(freq=f'{5}T',start='2020-10-10',periods=(12)*24*5)) df['col'] = np.random.randint(1, 100, size= df.shape[0]) df['uid'] = 1 df2 = pd.DataFrame(index=pd.date_range(freq=f'{5}T',start='2020-10-10',periods=(12)*24*5)) df2['col'] = np.random.randint(1, 50, size= df2.shape[0]) df2['uid'] = 2 df3=pd.concat([df, df2]).reset_index().rename(columns={'index': 'timestamp'})
输出样例:
timestamp col uid 0 2020-10-10 00:00:00 96 1 1 2020-10-10 00:05:00 47 1 2 2020-10-10 00:10:00 78 1 3 2020-10-10 00:15:00 27 1 ...
日志数据
import datetime as dt df_log=pd.DataFrame(np.array([[100, 1, 3], [40, 2, 6], [50, 1, 5], [60, 2, 9], [20, 1, 2], [30, 2, 5]]), columns=['duration', 'uid', 'factor']) df_log['timestamp'] = pd.Series([dt.datetime(2020,10,10, 15,21), dt.datetime(2020,10,10, 16,27), dt.datetime(2020,10,11, 21,25), dt.datetime(2020,10,11, 10,12), dt.datetime(2020,10,13, 20,56), dt.datetime(2020,10,13, 13,15)])
输出样例:
duration uid factor timestamp 0 100 1 3 2020-10-10 15:21:00 1 40 2 6 2020-10-10 16:27:00 ...
需求描述
将两份数据合并为df_merged,按uid分组完成以下操作:
- 生成新列
new,计算规则:df_merged['new'] = df_merged['duration'] * df_merged['factor'] - 将
new值按uid向前填充至下一条日志出现的时间点 - 叠加后续日志的计算值,最终实现2天移动求和(即每个时间点的
new值为过去2天内所有生效日志的duration*factor之和)
预期输出样例:
timestamp col uid duration factor new 0 2020-10-10 15:20:00 96 1 100 3 300 1 2020-10-10 15:25:00 47 1 100 3 300 2 2020-10-10 15:30:00 78 1 100 3 300 ... 2020-10-11 21:25:00 .. 1 60 9 540+300 2020-10-11 21:30:00 .. 1 60 9 540+300 ... 2020-10-13 20:55:00 .. 1 20 2 40+540 2020-10-13 21:00:00 .. 1 20 2 40+540 .. 2020-10-13 21:25:00 .. 1 20 2 40
实现思路
核心逻辑
每个日志的new值仅在日志时间戳到日志时间戳+2天的范围内生效,因此可以通过事件驱动的方式,标记每个日志的生效起始和结束节点,再通过累积和计算每个时间点的总和。
分步操作
预处理日志数据:
- 计算每条日志的
new值:df_log['new'] = df_log['duration'] * df_log['factor'] - 标记每条日志的生效结束时间:
df_log['end_time'] = df_log['timestamp'] + pd.Timedelta(days=2)
- 计算每条日志的
生成事件表:
- 将每条日志拆分为两个事件:在日志时间点添加
new值,在结束时间点减去new值 - 将事件的时间戳对齐到时间序列的5分钟粒度(向下取整)
- 将每条日志拆分为两个事件:在日志时间点添加
合并与累积求和:
- 将事件表与时间序列数据按
uid和timestamp左连接,缺失的事件值填充为0 - 按
uid和时间排序后,计算delta列的累积和,得到每个时间点的2天移动求和值
- 将事件表与时间序列数据按
格式补全(可选):
- 将原日志的
duration和factor字段合并到结果中,并按uid向前填充,匹配预期输出格式
- 将原日志的
完整代码示例
import pandas as pd import numpy as np import datetime as dt # 生成时间序列数据 df = pd.DataFrame(index=pd.date_range(freq=f'{5}T',start='2020-10-10',periods=(12)*24*5)) df['col'] = np.random.randint(1, 100, size= df.shape[0]) df['uid'] = 1 df2 = pd.DataFrame(index=pd.date_range(freq=f'{5}T',start='2020-10-10',periods=(12)*24*5)) df2['col'] = np.random.randint(1, 50, size= df2.shape[0]) df2['uid'] = 2 df3=pd.concat([df, df2]).reset_index().rename(columns={'index': 'timestamp'}) # 生成日志数据 df_log=pd.DataFrame(np.array([[100, 1, 3], [40, 2, 6], [50, 1, 5], [60, 2, 9], [20, 1, 2], [30, 2, 5]]), columns=['duration', 'uid', 'factor']) df_log['timestamp'] = pd.Series([dt.datetime(2020,10,10, 15,21), dt.datetime(2020,10,10, 16,27), dt.datetime(2020,10,11, 21,25), dt.datetime(2020,10,11, 10,12), dt.datetime(2020,10,13, 20,56), dt.datetime(2020,10,13, 13,15)]) # 预处理日志:计算new值和生效结束时间 df_log['new'] = df_log['duration'] * df_log['factor'] df_log['end_time'] = df_log['timestamp'] + pd.Timedelta(days=2) # 生成事件表 events = [] for _, row in df_log.iterrows(): # 生效起始事件:添加new值 events.append({'timestamp': row['timestamp'], 'uid': row['uid'], 'delta': row['new']}) # 生效结束事件:减去new值 events.append({'timestamp': row['end_time'], 'uid': row['uid'], 'delta': -row['new']}) events_df = pd.DataFrame(events) # 对齐事件时间到5分钟粒度 events_df['timestamp'] = events_df['timestamp'].dt.floor('5T') # 合并事件与时间序列数据 df_merged = pd.merge(df3, events_df, on=['uid', 'timestamp'], how='left') df_merged['delta'] = df_merged['delta'].fillna(0) # 按uid分组计算累积和,得到最终new列 df_merged = df_merged.sort_values(['uid', 'timestamp']) df_merged['new'] = df_merged.groupby('uid')['delta'].cumsum() # 合并原日志的duration和factor并向前填充(匹配预期格式) df_merged = pd.merge(df_merged, df_log[['uid', 'timestamp', 'duration', 'factor']], on=['uid', 'timestamp'], how='left') df_merged[['duration', 'factor']] = df_merged.groupby('uid')[['duration', 'factor']].ffill() # 查看结果 print(df_merged.loc[df_merged['uid'] == 1].head(30))
内容的提问来源于stack exchange,提问作者prof32
相关产品推荐
相关产品推荐

