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

Spark DataFrame中基于行值序列的复杂行分组实现(无UDF)

Spark实现事件子分组(无UDF方案)

问题背景

现有一组以EventId唯一标识的事件行,归属于GroupId标识的组。其中:

  • BeginEndMarker=1:起始事件
  • BeginEndMarker=5:结束事件
  • BeginEndMarker=-1:中间事件

需按规则生成子分组:

  • 每个子组以**首次出现的起始事件(BeginEndMarker=1)**为起始
  • 子组以上一个起始事件之前的最后一个结束事件(BeginEndMarker=5)为结束(允许无结束事件的不完整组)
  • 连续的起始事件归属于同一个子组
  • 要求不使用UDF,通过Spark原生API实现。

示例输入DataFrame

val df= Seq(
("GroupId1", "WF1", 1, "01-01-2023"),
("GroupId1", "WF2", -1, "01-02-2023"),
("GroupId1", "WF3", -1, "01-03-2023"),
("GroupId1", "WF4", 5, "01-04-2023"),
("GroupId1", "WF5", 5, "01-05-2023"),
("GroupId1", "WF6", 1, "01-06-2023"),
("GroupId1", "WF7", 1, "01-06-2023"),
("GroupId1", "WF8", -1, "01-07-2023"),
("GroupId1", "WF9", 5, "01-08-2023"),
("GroupId1", "WF10", 1, "01-09-2023"),
("GroupId1", "WF11", -1, "01-10-2023"),
).toDF("GroupId", "EventId","BeginEndMarker","Time")
df.show()

期望输出结果

+--------+-------+--------------+----------+--------+
| GroupId|EventId|BeginEndMarker|      Time|Subgroup|
+--------+-------+--------------+----------+--------+
|GroupId1|    WF1|             1|01-01-2023|     SG1|
|GroupId1|    WF2|            -1|01-02-2023|     SG1|
|GroupId1|    WF3|            -1|01-03-2023|     SG1|
|GroupId1|    WF4|             5|01-04-2023|     SG1|
|GroupId1|    WF5|             5|01-05-2023|     SG1|
|GroupId1|    WF6|             1|01-06-2023|     SG2|
|GroupId1|    WF7|             1|01-06-2023|     SG2|
|GroupId1|    WF8|            -1|01-07-2023|     SG2|
|GroupId1|    WF9|             5|01-08-2023|     SG2|
|GroupId1|   WF10|             1|01-09-2023|     SG3|
|GroupId1|   WF11|            -1|01-10-2023|     SG3|
+--------+-------+--------------+----------+--------+

实现方案

核心思路

通过窗口函数识别新子组触发点,再累计求和生成子组编号:

  1. 按GroupId分区、Time排序,保证同一组内事件按时间顺序处理;
  2. 用lag函数判断当前行是否为新子组起始:
    • 当前行是起始事件(BeginEndMarker=1)
    • 且上一行不是起始事件(或为分组内第一行)
  3. 对触发点标记为1,其余标记为0,累计求和得到子组编号;
  4. 拼接编号生成SGX格式的子组名称。

代码实现

// 导入Spark函数和窗口类
import org.apache.spark.sql.functions._
import org.apache.spark.sql.expressions.Window

// 将字符串时间转为日期类型,避免字符串排序异常
val dfWithTime = df.withColumn("Time", to_date(col("Time"), "dd-MM-yyyy"))

// 定义窗口:按GroupId分区,按Time升序排序
val windowSpec = Window.partitionBy("GroupId").orderBy("Time")

// 生成Subgroup列
val resultDf = dfWithTime
  // 标记新子组的触发点
  .withColumn("is_new_subgroup", when(
    col("BeginEndMarker") === 1 && 
    (lag(col("BeginEndMarker"), 1).over(windowSpec).isNull || lag(col("BeginEndMarker"), 1).over(windowSpec) =!= 1),
    1
  ).otherwise(0))
  // 累计求和得到子组编号
  .withColumn("subgroup_num", sum(col("is_new_subgroup")).over(windowSpec))
  // 拼接成SGX格式的子组名称
  .withColumn("Subgroup", concat(lit("SG"), col("subgroup_num")))
  // 清理中间临时列
  .drop("is_new_subgroup", "subgroup_num")

// 展示最终结果
resultDf.show()

逻辑验证

  • WF1:分组内第一行且是起始事件,触发新子组,编号为1 → SG1
  • WF6:上一行是结束事件(5),当前是起始事件,触发新子组,编号为2 → SG2
  • WF7:上一行是起始事件(1),不触发新子组,编号保持2 → SG2
  • WF10:上一行是结束事件(5),当前是起始事件,触发新子组,编号为3 → SG3
  • 所有中间事件、结束事件均继承当前子组编号,完全符合预期结果。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.21 23:27:14