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

如何基于Spark中前一值分组并计算机器运行时长?

问题描述

我有某台机器的运行数据,机器运行时每5秒至少生成一条含timestamp字段的记录,需要计算该机器每次开机的运行时长,即每次开机时段的首尾记录时间间隔。

我的初始思路是:

  • 先按timestamp排序数据集
  • 结合当前值与前一值(无前置值时用零值)生成timestamp_start和timestamp_now列:
    • 若当前timestamp与前一条的timestamp_now间隔超过5秒,则两列均设为当前timestamp
    • 若间隔≤5秒,则timestamp_start沿用前一条的值,timestamp_now设为当前timestamp
  • 之后按timestamp_start分组取最大timestamp_now,再计算时长

但我不确定该用fold、agg还是reduce,也考虑过滑动窗口,缺乏相关经验。作为Spark新手,我尝试了以下代码:

初始代码尝试

DataQuery.builder(spark).variables() \
    .system('XXX') \
    .nameLike('XXX%XXX%') \
    .timeWindow('2021-10-10 00:00:00.000', '2022-11-28 00:00:00.000') \
    .build() \
    .orderBy('timestamp')
    .agg('timestamp', # How do I get to the previous entry?)

后续尝试代码

df = DataQuery.builder(spark).variables() \
    .system('XXX') \
    .nameLike('XXX') \
    .timeWindow('2021-08-10 00:00:00.000', '2022-11-28 00:00:00.000') \
    .build()

timestamps = df.sort('timestamp') \
    .select(psf.from_unixtime('nxcals_timestamp').alias('ts'))

# AT LEAST I HOPE THIS LINE IS RIGHT (?)
window = timestamps.groupBy(psf.session_window('ts', '10 minutes')) \
    .agg(psf.min(timestamps.ts))

window_timestamps = window.select(window.session_window.start.cast("string").alias("start"), window.session_window.end.cast("string").alias('end'))

执行show()后得到的结果:

+--------------------+--------------------+
|               start|                 end|
+--------------------+--------------------+
|-290308-12-21 20:...|-290308-12-21 20:...|
|-290308-12-23 17:...|-290308-12-23 17:...|
|-290308-12-25 06:...|-290308-12-25 06:...|
|-290308-12-25 15:...|-290308-12-25 15:...|
|-290307-01-01 05:...|-290307-01-01 05:...|
|-290307-01-04 06:...|-290307-01-04 06:...|
|-290307-01-04 19:...|-290307-01-04 19:...|
|-290307-01-05 05:...|-290307-01-05 05:...|
|-290307-01-05 08:...|-290307-01-05 08:...|
|-290307-01-06 00:...|-290307-01-06 00:...|
|-290307-01-10 07:...|-290307-01-10 07:...|
|-290307-01-14 11:...|-290307-01-14 11:...|
|-290307-01-15 03:...|-290307-01-15 04:...|
|-290307-01-15 08:...|-290307-01-15 08:...|
|-290307-01-15 13:...|-290307-01-15 13:...|
|-290307-01-16 17:...|-290307-01-16 17:...|
|-290307-01-20 16:...|-290307-01-20 16:...|
|-290307-01-24 19:...|-290307-01-24 19:...|
|-290307-01-26 17:...|-290307-01-26 17:...|
|-290307-01-30 23:...|-290307-01-30 23:...|
+--------------------+--------------------+

现在我需要将这些数据转换为时间差列,尝试以下代码但无法正常运行:

diff = window_timestamps.rdd.map(lambda row: row.end.cast('long') - row.start.cast('long')).toDF(["diff_in_seconds"])

解决方案

核心问题分析

你当前的代码存在两个关键问题:

  1. 会话窗口超时参数设置错误(10分钟不符合“5秒内有记录即视为同一段开机时间”的业务逻辑)
  2. 时间格式转换错误,导致出现远古时间值,后续时间差计算完全失效

正确实现步骤

1. 修复时间格式转换

首先确保nxcals_timestamp转换为正确的Timestamp类型,需注意时间戳单位(毫秒/秒):

# 假设nxcals_timestamp是毫秒级数值型时间戳,转换为Spark Timestamp类型
timestamps = df.sort('nxcals_timestamp') \
    .select(psf.to_timestamp(psf.col('nxcals_timestamp')/1000).alias('ts'))

2. 用会话窗口划分开机时段

将会话窗口的超时时间设为5秒,超过5秒无记录则视为机器停机:

# 定义会话窗口:连续记录间隔≤5秒归为同一会话(开机时段)
session_window_spec = psf.session_window('ts', '5 seconds')
window_df = timestamps.groupBy(session_window_spec) \
    .agg(
        psf.min('ts').alias('start_time'),
        psf.max('ts').alias('end_time')
    )

3. 计算运行时长

直接用Spark内置函数计算时间差,无需转RDD:

# 计算每个开机时段的运行时长(单位:秒)
result_df = window_df.select(
    psf.col('session_window.start').alias('session_start'),
    psf.col('session_window.end').alias('session_end'),
    psf.col('start_time'),
    psf.col('end_time'),
    (psf.unix_timestamp('end_time') - psf.unix_timestamp('start_time')).alias('run_duration_seconds')
)

result_df.show()

4. 替代方案:用窗口函数实现初始思路

如果不想用会话窗口,也可以通过窗口函数标记开机时段:

# 定义排序窗口,用于获取前一条记录的时间
lag_window = psf.Window.orderBy('ts')

# 计算当前记录与前一条的时间差
timestamps_with_diff = timestamps.withColumn(
    'prev_ts', psf.lag('ts').over(lag_window)
).withColumn(
    'time_diff', psf.unix_timestamp('ts') - psf.unix_timestamp('prev_ts')
)

# 标记新的开机时段:第一条记录或时间差>5秒时,session_id自增1
timestamps_with_session = timestamps_with_diff.withColumn(
    'session_id', psf.sum(
        psf.when(psf.col('time_diff') > 5 | psf.col('prev_ts').isNull(), 1).otherwise(0)
    ).over(lag_window)
)

# 按session_id分组计算每个时段的首尾时间和时长
final_result = timestamps_with_session.groupBy('session_id') \
    .agg(
        psf.min('ts').alias('start_time'),
        psf.max('ts').alias('end_time')
    ).withColumn(
        'run_duration_seconds', psf.unix_timestamp('end_time') - psf.unix_timestamp('start_time')
    )

错误代码修复

针对你最后那段无法运行的代码,问题在于字符串类型不能直接转long,需先转回Timestamp类型再计算时间差:

diff = window_timestamps.select(
    (psf.unix_timestamp(psf.col('end')) - psf.unix_timestamp(psf.col('start'))).alias('diff_in_seconds')
)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.07 13:00:46