Spark数据集分组求助:基于当前与前序状态的分组逻辑实现
问题解决:按状态边界生成分组ID
你写的代码只是按日期生成了行号,完全没用到前后状态的对比逻辑,所以肯定实现不了需求。要实现你要的分组逻辑,核心是检测到「前一条是F、当前是C」的边界时,分组ID加1,具体可以用Spark的窗口函数lag(取前一行数据)和累加求和来实现,代码如下:
import org.apache.spark.sql.expressions.Window; import org.apache.spark.sql.expressions.WindowSpec; import static org.apache.spark.sql.functions.*; // 定义窗口:按Date排序,确保数据顺序正确 WindowSpec orderWindow = Window.orderBy(col("Date")); Dataset<Row> resultDs = dataset // 获取前一条记录的status,第一条记录没有前一条,用null填充 .withColumn("prev_status", lag(col("status"), 1).over(orderWindow)) // 标记分组边界:当前是C且前一条是F时,记为1,否则0 .withColumn("flag", when((col("status").equalTo("C").and(col("prev_status").equalTo("F"))), 1).otherwise(0)) // 累加flag值,再+1得到分组ID(因为初始累加从0开始,加1后从1计数) .withColumn("id", sum(col("flag")).over(orderWindow.rowsBetween(Window.unboundedPreceding, Window.currentRow)).plus(1)) // 去掉中间临时列 .drop("prev_status", "flag") .persist(StorageLevel.MEMORY_ONLY());
代码逻辑说明:
lag函数:用来获取当前行的前一行status,保证我们能对比前后状态。flag列:精准捕捉你要的分组边界——只有当前是C且前一条是F时,才触发新分组,标记为1。- 累加求和生成ID:对
flag列从第一行到当前行做累加,每遇到一个1,累加值就加1,最后加1让分组ID从1开始,正好对应你要的分组结果。
运行这段代码后,输出就会和你期望的完全一致:
| Date | status | id |
|---|---|---|
| 20230101 | C | 1 |
| 20230101 | C | 1 |
| 20230101 | F | 1 |
| 20230101 | C | 2 |
| 20230101 | C | 2 |
| 20230101 | F | 2 |
| 20230101 | C | 3 |
| 20230101 | F | 3 |
| 20230101 | F | 3 |
内容的提问来源于stack exchange,提问作者user10146200
相关产品推荐
相关产品推荐

