如何在Pandas中按动态时间窗口聚合短时间发送的消息?
需求
将发送时间间隔在指定时长内的消息分组为一个块,若有消息加入块,时间窗口将延长以纳入更多符合条件的消息。
示例输入
| datetime | message | |
|---|---|---|
| 0 | 2023-01-01 12:00:00 | A |
| 1 | 2023-01-01 12:20:00 | B |
| 2 | 2023-01-01 12:30:00 | C |
| 3 | 2023-01-01 12:30:55 | D |
| 4 | 2023-01-01 12:31:20 | E |
| 5 | 2023-01-01 15:00:00 | F |
| 6 | 2023-01-01 15:30:30 | G |
| 7 | 2023-01-01 15:30:55 | H |
参数设为1分钟时的预期输出
| datetime | message | datetime_last | n_block | |
|---|---|---|---|---|
| 0 | 2023-01-01 12:00:00 | A | 2023-01-01 12:00:00 | 1 |
| 1 | 2023-01-01 12:20:00 | B | 2023-01-01 12:20:00 | 1 |
| 2 | 2023-01-01 12:30:00 | C\nD\nE | 2023-01-01 12:31:20 | 3 |
| 3 | 2023-01-01 15:00:00 | F | 2023-01-01 15:00:00 | 1 |
| 4 | 2023-01-01 15:30:30 | G\nH | 2023-01-01 15:30:55 | 2 |
失败尝试
尝试用rolling窗口实现消息连续拼接,代码如下:
def join_messages(x): return '\n'.join(x) df.rolling(window='1min', on='datetime').agg({ 'datetime': ['first', 'last'], 'message': [join_messages, "count"]}) # 尝试用聚合后的datetime.first覆盖原datetime列
运行后触发ValueError: invalid on specified as datetime, must be a column (of DataFrame), an Index or None,无法在窗口中正常调用datetime列,且rolling对字符串处理效果差,此方案不可行,需更简洁的解决方法。
输入与预期输出的代码片段
import pandas as pd df = pd.DataFrame({ 'datetime': [pd.Timestamp('2023-01-01 12:00'), pd.Timestamp('2023-01-01 12:20'), pd.Timestamp('2023-01-01 12:30:00'), pd.Timestamp('2023-01-01 12:30:55'), pd.Timestamp('2023-01-01 12:31:20'), pd.Timestamp('2023-01-01 15:00'), pd.Timestamp('2023-01-01 15:30:30'), pd.Timestamp('2023-01-01 15:30:55'),], 'message': list('ABCDEFGH')}) df_expected = pd.DataFrame({ 'datetime': [pd.Timestamp('2023-01-01 12:00'), pd.Timestamp('2023-01-01 12:20'), pd.Timestamp('2023-01-01 12:30:00'), pd.Timestamp('2023-01-01 15:00'), pd.Timestamp('2023-01-01 15:30:30'),], 'message': ['A', 'B', 'C\nD\nE', 'F', 'G\nH'], 'datetime_last': [pd.Timestamp('2023-01-01 12:00'), pd.Timestamp('2023-01-01 12:20'), pd.Timestamp('2023-01-01 12:31:20'), pd.Timestamp('2023-01-01 15:00'), pd.Timestamp('2023-01-01 15:30:55'),], 'n_block': [1, 1, 3, 1, 2]})
解决方案
rolling窗口为固定范围,不适合这种动态延长窗口的需求,以下两种方法可实现目标:
方法1:迭代标记分组(逻辑直观)
def group_messages(df, window_minutes=1): window = pd.Timedelta(minutes=window_minutes) groups = [] current_group = [df.iloc[0]] for idx in range(1, len(df)): row = df.iloc[idx] # 当前消息与组内最后一条消息的间隔在窗口内则加入当前组 if row['datetime'] - current_group[-1]['datetime'] <= window: current_group.append(row) else: groups.append(current_group) current_group = [row] # 加入最后一组 groups.append(current_group) # 聚合每组数据 result = [] for group in groups: group_df = pd.DataFrame(group) result.append({ 'datetime': group_df['datetime'].iloc[0], 'message': '\n'.join(group_df['message']), 'datetime_last': group_df['datetime'].iloc[-1], 'n_block': len(group_df) }) return pd.DataFrame(result) # 调用函数 df_result = group_messages(df, window_minutes=1) # 验证与预期一致(无输出则匹配) print(pd.testing.assert_frame_equal(df_result, df_expected))
方法2:向量化标记分组(高效适配大数据量)
def group_messages_vectorized(df, window_minutes=1): window = pd.Timedelta(minutes=window_minutes) # 计算每条消息与前一条的时间差 diffs = df['datetime'].diff().fillna(pd.Timedelta(0)) # 标记新分组起始点:时间差超过窗口则为新组 group_ids = (diffs > window).cumsum() # 按分组聚合 result = df.groupby(group_ids).agg( datetime=('datetime', 'first'), message=('message', '\n'.join), datetime_last=('datetime', 'last'), n_block=('message', 'count') ).reset_index(drop=True) return result # 调用函数 df_result = group_messages_vectorized(df, window_minutes=1) # 验证与预期一致(无输出则匹配) print(pd.testing.assert_frame_equal(df_result, df_expected))
说明
- 两种方法均实现动态窗口逻辑:只要新消息与当前组最后一条消息的间隔在指定时长内,就加入当前组,窗口自动延长
- 向量化方法利用
groupby聚合,效率远高于迭代,适合处理大规模数据
内容的提问来源于stack exchange,提问作者Fabitosh
相关产品推荐
相关产品推荐

