如何基于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"])
解决方案
核心问题分析
你当前的代码存在两个关键问题:
- 会话窗口超时参数设置错误(10分钟不符合“5秒内有记录即视为同一段开机时间”的业务逻辑)
- 时间格式转换错误,导致出现远古时间值,后续时间差计算完全失效
正确实现步骤
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
相关产品推荐
相关产品推荐

