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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 19:18:19