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

使用PySpark根据连接断开时间戳统计每日活跃设备数

PySpark 按日统计IoT活跃设备实现方案

核心判断规则:只要某日期处于设备的连接时间~断开时间区间内,就算该设备当日活跃,跨天连接的设备需要在覆盖的所有日期都计入统计。


实现步骤说明

  1. 提取时间戳里的日期值,将连接时间、断开时间统一转为Date类型,忽略时分秒精度
  2. 为每条连接记录生成连接日期到断开日期之间的所有连续日期,每个日期对应当日该设备活跃
  3. 按日期分组,统计去重后的设备总数,得到当日活跃设备量,避免同一设备同一天多次连接被重复计数
  4. 格式化日期输出为dd-MM-yyyy格式,匹配预期输出样式

完整实现代码

from pyspark.sql import SparkSession
from pyspark.sql import functions as F

# 初始化Spark会话(已有现成会话可跳过这步)
spark = SparkSession.builder.appName("iot_active_device_stat").getOrCreate()

# 假设原始DataFrame名为df,字段为device_id、connected_at、disconnect_at
result_df = df.withColumn(
    # 连接时间转日期
    "conn_date", 
    F.to_date("connected_at")
).withColumn(
    # 断开时间转日期,若断开时间为空默认取当前日期(适配设备持续在线未断开场景)
    "disc_date",
    F.coalesce(F.to_date("disconnect_at"), F.current_date())
).withColumn(
    # 生成连接到断开之间的所有日期序列,炸开为单行记录
    "date",
    F.explode(F.sequence("conn_date", "disc_date", F.expr("interval 1 day")))
).groupBy("date").agg(
    # 统计去重后的设备数,命名为active users
    F.countDistinct("device_id").alias("active users")
).withColumn(
    # 日期格式化为dd-MM-yyyy字符串
    "date",
    F.date_format("date", "dd-MM-yyyy")
).orderBy("date")

# 打印查看结果
result_df.show()

输出示例

执行后输出格式和预期完全对齐:

+----------+------------+
|      date|active users|
+----------+------------+
|20-08-2021|           1|
+----------+------------+

注意事项

  • 可提前过滤脏数据:比如断开时间早于连接时间的异常记录、连接时长过短的无效测试记录,减少后续计算的数据量
  • 跨天连接记录炸开时会产生数据膨胀,如果数据规模极大,建议先做数据清洗再执行序列生成逻辑
  • 如果需要统计历史截止到某一天的活跃数,可以把disc_date的默认值替换为指定的统计截止日期即可

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 21:54:34