基于PySpark计算跨多天时间戳的单设备每日可用时长
PySpark计算设备每日可用时长(无需UDF)
核心思路
完全依赖PySpark内置日期/时间函数,将跨天故障事件拆分为单日记录,计算每日停机时长后推导可用时长与可用率,最后补全目标日期范围内的所有设备-日期组合,全程无需自定义UDF。
具体实现步骤
1. 数据预处理:转换日期类型
先将startDate和endDate转换为Timestamp类型,确保后续时间计算的准确性:
from pyspark.sql import functions as F from pyspark.sql.types import TimestampType # 假设原始数据集为df df = df.withColumn("startDate", F.to_timestamp("startDate")) \ .withColumn("endDate", F.to_timestamp("endDate"))
2. 拆分跨天事件为单日记录
用sequence函数生成故障事件覆盖的所有日期数组,再通过explode展开为单行单日记录:
df_days = df.withColumn("event_days", F.sequence( F.date_trunc("day", "startDate"), F.date_trunc("day", "endDate"), F.expr("INTERVAL 1 DAY") )) \ .explode("event_days") \ .withColumn("day", F.to_date("event_days"))
3. 计算单日停机时长
对每个单日记录,取当天内实际的停机区间,再转换为秒数:
df_downtime = df_days.withColumn( "downtime_start", F.greatest("startDate", "event_days") ).withColumn( "downtime_end", F.least("endDate", F.date_add("event_days", 1)) ).withColumn( "downtime_seconds", F.unix_timestamp("downtime_end") - F.unix_timestamp("downtime_start") )
说明:event_days为当天0点,date_add(event_days,1)为次日0点,通过greatest和least锁定当天内的停机时间段,再计算时间差得到单日停机秒数。
4. 聚合每日总停机时长
按deviceId和day分组,求和单日所有故障的总停机时长:
df_daily_downtime = df_downtime.groupBy("deviceId", "day") \ .agg(F.sum("downtime_seconds").alias("total_downtime"))
5. 计算可用时长与可用率
以一天总秒数(86400秒)为基准,推导可用时长和可用率:
df_uptime = df_daily_downtime.withColumn( "uptime_seconds", 86400 - F.col("total_downtime") ).withColumn( "uptime_rate", F.round((F.col("uptime_seconds") / 86400) * 100, 1) ).withColumn( "uptime", F.concat(F.col("uptime_rate"), F.lit("%")) ).select("deviceId", "day", "uptime")
6. 补全目标日期范围内的所有记录
若需要覆盖指定日期范围(如示例中的2022-06-10至2022-08-08),生成设备与日期的笛卡尔积后左连接补全无故障日期:
# 生成目标日期序列 start_date = "2022-06-10" end_date = "2022-08-08" dates_df = spark.createDataFrame( [(d,) for d in spark.sql(f"select sequence(to_date('{start_date}'), to_date('{end_date}'), interval 1 day) as dates").first()[0]], ["day"] ) # 获取所有唯一设备 devices_df = df.select("deviceId").distinct() # 生成设备-日期笛卡尔积 all_device_days = devices_df.crossJoin(dates_df) # 左连接补全数据,无故障日期默认可用率100% final_df = all_device_days.join(df_uptime, on=["deviceId", "day"], how="left") \ .withColumn("uptime", F.coalesce("uptime", F.lit("100%")))
结果示例
最终final_df输出格式与预期一致:
+--------+----------+-------+ |deviceId|day |uptime | +--------+----------+-------+ |11615 |2022-06-10|100% | |11615 |2022-06-11|88.3% | |11615 |2022-06-12|0% | |11615 |2022-06-13|76.0% | |11618 |2022-06-10|100% | +--------+----------+-------+
内容的提问来源于stack exchange,提问作者Tristan Foureur
相关产品推荐
相关产品推荐

