PySpark中按cd_equipment_no分组并为批次生成独立序号
问题
需要分析同一设备编号(cd_equipment_no)在不同时间批次内的组合数量。现有PySpark数据集包含ID、source_timestamp(已转为时间戳类型)、node_id、cd_equipment_no等字段。
尝试过用窗口函数按cd_equipment_no分组、source_timestamp升序排序生成row_number,但得到的是设备组内的全局连续序号,不符合需求。实际需要同一cd_equipment_no下的不同时间批次内序号从1开始重新计数。
PySpark实现方案
实现核心是先给同一设备的记录划分时间批次,再在每个批次内生成序号。根据批次定义的不同,分两种常见场景实现:
场景1:按固定时间间隔划分批次(如小时/天/分钟)
适用于批次有固定时间周期的场景,比如每小时一个批次:
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window # 初始化SparkSession spark = SparkSession.builder.appName("EquipmentFixedBatch").getOrCreate() # 假设已有数据集df,包含需求中的所有字段 # 1. 生成固定间隔的批次ID,这里以小时为单位,可替换为"day"/"minute"等 df_with_batch = df.withColumn( "batch_id", F.date_trunc("hour", F.col("source_timestamp")) ) # 2. 定义窗口:按设备编号+批次分组,按时间升序排序 window_spec = Window.partitionBy("cd_equipment_no", "batch_id").orderBy("source_timestamp") # 3. 生成批次内的序号 result_df = df_with_batch.withColumn( "batch_row_num", F.row_number().over(window_spec) ).select("ID", "source_timestamp", "node_id", "cd_equipment_no", "batch_id", "batch_row_num") # 查看结果 result_df.show()
场景2:按设备活跃间隔划分批次(如超过N分钟无数据则视为新批次)
适用于批次由设备活跃中断触发的场景,比如设备连续30分钟无数据后重启视为新批次:
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window spark = SparkSession.builder.appName("EquipmentDynamicBatch").getOrCreate() # 1. 按设备分组排序,计算当前记录与上一条的时间差(分钟) window_order = Window.partitionBy("cd_equipment_no").orderBy("source_timestamp") df_with_diff = df.withColumn( "prev_timestamp", F.lag("source_timestamp").over(window_order) ).withColumn( "time_diff_minutes", F.when( F.col("prev_timestamp").isNotNull(), (F.col("source_timestamp").cast("long") - F.col("prev_timestamp").cast("long")) / 60 ).otherwise(0) ) # 2. 标记新批次:时间差超过30分钟则标记为新批次起始点 df_with_batch_flag = df_with_diff.withColumn( "is_new_batch", F.when(F.col("time_diff_minutes") > 30, 1).otherwise(0) ) # 3. 累加标记生成批次ID df_with_batch = df_with_batch_flag.withColumn( "batch_id", F.sum("is_new_batch").over(window_order.rangeBetween(Window.unboundedPreceding, 0)) ) # 4. 生成批次内的序号 window_batch = Window.partitionBy("cd_equipment_no", "batch_id").orderBy("source_timestamp") result_df = df_with_batch.withColumn( "batch_row_num", F.row_number().over(window_batch) ).select("ID", "source_timestamp", "node_id", "cd_equipment_no", "batch_id", "batch_row_num") result_df.show()
说明
- 固定间隔场景可通过修改
date_trunc的参数调整批次周期; - 动态间隔场景可修改
time_diff_minutes > 30中的阈值,适配实际业务的活跃中断定义; - 最终的
batch_row_num即为每个批次内从1开始的序号,结合cd_equipment_no和batch_id即可统计各批次的组合数量。
内容的提问来源于stack exchange,提问作者Gaaaa
相关产品推荐
相关产品推荐

