如何使用Spark DataFrame基于状态计算时间戳累计差值?
基于ID和STATUS计算连续行时间戳累计差值
需要针对DataFrame中的数据,按ID和连续的STATUS分组,计算每行时间戳与组内首行时间戳的累计差值(单位:秒),组内首行差值为0。
示例输入DataFrame
| ID | TIMESTAMP | STATUS |
|---|---|---|
| V1 | 2023-06-18 13:00:00 | 1 |
| V1 | 2023-06-18 13:01:00 | 1 |
| V1 | 2023-06-18 13:02:00 | 1 |
| V1 | 2023-06-18 13:03:00 | 1 |
| V1 | 2023-06-18 13:06:00 | 0 |
| V1 | 2023-06-18 13:07:00 | 0 |
| V1 | 2023-06-18 13:08:00 | 0 |
| V1 | 2023-06-18 13:09:00 | 1 |
| V1 | 2023-06-18 13:10:00 | 1 |
| V1 | 2023-06-18 13:11:00 | 1 |
| V1 | 2023-06-18 13:12:00 | 1 |
预期输出结果
| ID | TIMESTAMP | STATUS | TIMESTAMP_DIFF(in_seconds) |
|---|---|---|---|
| V1 | 2023-06-18 13:00:00 | 1 | 0 |
| V1 | 2023-06-18 13:01:00 | 1 | 60 |
| V1 | 2023-06-18 13:02:00 | 1 | 120 |
| V1 | 2023-06-18 13:03:00 | 1 | 180 |
| V1 | 2023-06-18 13:06:00 | 0 | 0 |
| V1 | 2023-06-18 13:07:00 | 0 | 60 |
| V1 | 2023-06-18 13:08:00 | 0 | 120 |
| V1 | 2023-06-18 13:09:00 | 1 | 0 |
| V1 | 2023-06-18 13:10:00 | 1 | 60 |
| V1 | 2023-06-18 13:11:00 | 1 | 120 |
| V1 | 2023-06-18 13:12:00 | 1 | 180 |
实现方案(Python Pandas)
步骤说明:
- 将
TIMESTAMP列转换为datetime类型,确保能进行时间运算 - 生成连续状态分组标识:通过比较当前行
STATUS与上一行是否不同,结合ID分组,标记出每个连续相同ID+STATUS的组 - 按分组标识分组,计算每行
TIMESTAMP与组内首个TIMESTAMP的差值,转换为秒数
代码示例:
import pandas as pd # 构造示例数据 data = { 'ID': ['V1']*11, 'TIMESTAMP': [ '2023-06-18 13:00:00', '2023-06-18 13:01:00', '2023-06-18 13:02:00', '2023-06-18 13:03:00', '2023-06-18 13:06:00', '2023-06-18 13:07:00', '2023-06-18 13:08:00', '2023-06-18 13:09:00', '2023-06-18 13:10:00', '2023-06-18 13:11:00', '2023-06-18 13:12:00' ], 'STATUS': [1,1,1,1,0,0,0,1,1,1,1] } df = pd.DataFrame(data) # 1. 转换时间列为datetime类型 df['TIMESTAMP'] = pd.to_datetime(df['TIMESTAMP']) # 2. 生成连续状态分组标识:同一ID下,STATUS变化时生成新组 df['group_id'] = (df['STATUS'] != df['STATUS'].shift(1)) | (df['ID'] != df['ID'].shift(1)) df['group_id'] = df.groupby('ID')['group_id'].cumsum() # 3. 计算累计时间差(秒) df['TIMESTAMP_DIFF(in_seconds)'] = df.groupby(['ID', 'group_id'])['TIMESTAMP'].transform( lambda x: (x - x.iloc[0]).dt.total_seconds() ) # 可选:删除临时的group_id列 df = df.drop('group_id', axis=1) print(df)
代码解释:
df['TIMESTAMP'] = pd.to_datetime(df['TIMESTAMP']):将字符串格式的时间转换为Pandas可识别的datetime对象df['group_id'] = (df['STATUS'] != df['STATUS'].shift(1)) | (df['ID'] != df['ID'].shift(1)):判断当前行与上一行的ID或STATUS是否变化,变化则标记为新组的开始df.groupby('ID')['group_id'].cumsum():按ID分组后,对标记值累加,生成唯一的组IDgroupby(['ID', 'group_id'])['TIMESTAMP'].transform(...):按ID和组ID分组,计算每行时间与组内首行时间的差值并转成秒数
内容的提问来源于stack exchange,提问作者RMK
相关产品推荐
相关产品推荐

