You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用Spark DataFrame基于状态计算时间戳累计差值?

基于ID和STATUS计算连续行时间戳累计差值

需要针对DataFrame中的数据,按ID和连续的STATUS分组,计算每行时间戳与组内首行时间戳的累计差值(单位:秒),组内首行差值为0。

示例输入DataFrame

IDTIMESTAMPSTATUS
V12023-06-18 13:00:001
V12023-06-18 13:01:001
V12023-06-18 13:02:001
V12023-06-18 13:03:001
V12023-06-18 13:06:000
V12023-06-18 13:07:000
V12023-06-18 13:08:000
V12023-06-18 13:09:001
V12023-06-18 13:10:001
V12023-06-18 13:11:001
V12023-06-18 13:12:001

预期输出结果

IDTIMESTAMPSTATUSTIMESTAMP_DIFF(in_seconds)
V12023-06-18 13:00:0010
V12023-06-18 13:01:00160
V12023-06-18 13:02:001120
V12023-06-18 13:03:001180
V12023-06-18 13:06:0000
V12023-06-18 13:07:00060
V12023-06-18 13:08:000120
V12023-06-18 13:09:0010
V12023-06-18 13:10:00160
V12023-06-18 13:11:001120
V12023-06-18 13:12:001180

实现方案(Python Pandas)

步骤说明:

  1. 将TIMESTAMP列转换为datetime类型,确保能进行时间运算
  2. 生成连续状态分组标识:通过比较当前行STATUS与上一行是否不同,结合ID分组,标记出每个连续相同ID+STATUS的组
  3. 按分组标识分组,计算每行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分组后,对标记值累加,生成唯一的组ID
  • groupby(['ID', 'group_id'])['TIMESTAMP'].transform(...):按ID和组ID分组,计算每行时间与组内首行时间的差值并转成秒数

内容的提问来源于stack exchange,提问作者RMK

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.12 21:24:51