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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 21:45:33