Spark按序列分组记录求助:基于状态变化分组且从R开始
解决方案
要实现你需要的分组逻辑,核心是通过窗口函数对比当前行与前一行的状态,并根据状态是否相同来决定是否开启新分组。以下是具体实现步骤和代码:
逻辑梳理
根据你给出的输入和期望输出,分组规则可拆解为:
- 当当前行状态与前一行相同时,开启新分组(分组ID+1)
- 当当前行状态与前一行不同时,保持当前分组(分组ID不变)
- 分组ID从1开始,连续相同状态会触发分组递增,不同状态则归为同一组
实现代码
import org.apache.spark.sql.functions._ import org.apache.spark.sql.Window import org.apache.spark.storage.StorageLevel val resultDs = dataset // 添加行号确保数据顺序(因Date相同,必须用行号固定原始顺序) .withColumn("row_num", row_number().over(Window.orderBy("Date"))) // 获取前一行的status值 .withColumn("prev_status", lag("status", 1).over(Window.orderBy("row_num"))) // 标记是否需要递增分组ID:当前与前一行状态相同时标记为1,否则为0 .withColumn("incr", when(col("prev_status") === col("status"), 1).otherwise(0)) // 计算累计分组ID:初始为1,累加所有前置的标记值 .withColumn("id", 1 + sum("incr").over(Window.orderBy("row_num").rowsBetween(Window.unboundedPreceding, Window.currentRow))) // 保留目标列并持久化 .select("Date", "status", "id") .persist(StorageLevel.MEMORY_ONLY())
原代码问题说明
你之前的代码直接用row_number()作为分组ID,这只是按Date排序的行号,完全没有结合状态变化的逻辑,因此无法满足分组需求。
内容的提问来源于stack exchange,提问作者Tim-Timer
相关产品推荐
相关产品推荐

