如何基于连续时间行的相同值聚合PySpark DataFrame
需求描述
我有一个包含Time、Name、Flag三列的DataFrame,需要将Name和Flag值相同且时间连续的行聚合为Start和End列。这里的时间连续指相邻行的时间间隔恰好为1分钟(比如1:08和1:10因缺失1:09,不属于连续,不会合并)。
输入DataFrame
| Time | Name | Flag |
|---|---|---|
| 5/1/2023 1:01 | Peter | 1 |
| 5/1/2023 1:02 | Peter | 1 |
| 5/1/2023 1:03 | Peter | 1 |
| 5/1/2023 1:04 | Peter | 0 |
| 5/1/2023 1:05 | Peter | 0 |
| 5/1/2023 1:06 | Peter | 1 |
| 5/1/2023 1:07 | Peter | 1 |
| 5/1/2023 1:08 | Peter | 1 |
| 5/1/2023 1:01 | John | 1 |
| 5/1/2023 1:02 | John | 0 |
| 5/1/2023 1:03 | John | 0 |
| 5/1/2023 1:04 | John | 0 |
| 5/1/2023 1:05 | John | 0 |
| 5/1/2023 1:06 | John | 0 |
| 5/1/2023 1:07 | John | 1 |
| 5/1/2023 1:08 | John | 1 |
| 5/2/2023 1:10 | Peter | 1 |
| 5/2/2023 1:11 | Peter | 1 |
| 5/2/2023 1:20 | John | 0 |
| 5/2/2023 1:21 | John | 0 |
| 5/2/2023 1:22 | John | 0 |
期望输出DataFrame
| Start | End | Name | Flag |
|---|---|---|---|
| 5/1/2023 1:01 | 5/1/2023 1:03 | Peter | 1 |
| 5/1/2023 1:04 | 5/1/2023 1:05 | Peter | 0 |
| 5/1/2023 1:06 | 5/1/2023 1:08 | Peter | 1 |
| 5/2/2023 1:10 | 5/2/2023 1:11 | Peter | 1 |
| 5/1/2023 1:01 | 5/1/2023 1:01 | John | 1 |
| 5/1/2023 1:02 | 5/1/2023 1:06 | John | 0 |
| 5/1/2023 1:07 | 5/1/2023 1:08 | John | 1 |
| 5/2/2023 1:20 | 5/2/2023 1:22 | John | 0 |
解决方案
通过以下步骤用Pandas实现需求:
1. 转换时间列格式
先把Time列转为datetime类型,才能正确计算时间间隔:
import pandas as pd # 假设输入数据已加载到df变量中 df['Time'] = pd.to_datetime(df['Time'])
2. 标记连续时间段
按Name和Flag分组,判断当前行与上一行的时间差是否为1分钟,以此划分不同的连续组:
# 计算同组内相邻行的时间差(单位:分钟) df['time_diff'] = df.groupby(['Name', 'Flag'])['Time'].diff().dt.total_seconds() / 60 # 标记新组:时间差不等于1或为NaN时,视为新连续段的起点 df['group_id'] = (df['time_diff'] != 1).cumsum()
3. 聚合生成结果
按Name、Flag和group_id分组,取每组的最小时间作为Start,最大时间作为End,最后整理列顺序:
# 聚合分组结果 result = df.groupby(['Name', 'Flag', 'group_id']).agg( Start=('Time', 'min'), End=('Time', 'max') ).reset_index().drop(columns='group_id') # 调整列顺序为需求格式 result = result[['Start', 'End', 'Name', 'Flag']] # 可选:将时间格式转回原输入的字符串样式 result['Start'] = result['Start'].dt.strftime('%m/%d/%Y %H:%M') result['End'] = result['End'].dt.strftime('%m/%d/%Y %H:%M')
运行上述代码后,result即为符合要求的输出DataFrame。
内容的提问来源于stack exchange,提问作者n179911a
相关产品推荐
相关产品推荐

