使用PySpark根据连接断开时间戳统计每日活跃设备数
PySpark 按日统计IoT活跃设备实现方案
核心判断规则:只要某日期处于设备的连接时间~断开时间区间内,就算该设备当日活跃,跨天连接的设备需要在覆盖的所有日期都计入统计。
实现步骤说明
- 提取时间戳里的日期值,将连接时间、断开时间统一转为
Date类型,忽略时分秒精度 - 为每条连接记录生成连接日期到断开日期之间的所有连续日期,每个日期对应当日该设备活跃
- 按日期分组,统计去重后的设备总数,得到当日活跃设备量,避免同一设备同一天多次连接被重复计数
- 格式化日期输出为
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
相关产品推荐
相关产品推荐

