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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 11:00:26