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| +--------+-------+--------------+----------+--------+
实现方案
核心思路
通过窗口函数识别新子组触发点,再累计求和生成子组编号:
- 按
GroupId分区、Time排序,保证同一组内事件按时间顺序处理; - 用
lag函数判断当前行是否为新子组起始:- 当前行是起始事件(
BeginEndMarker=1) - 且上一行不是起始事件(或为分组内第一行)
- 当前行是起始事件(
- 对触发点标记为1,其余标记为0,累计求和得到子组编号;
- 拼接编号生成
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
相关产品推荐
相关产品推荐

