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

Apache Spark基于rec列检测键变更并生成groupId的实现问题

问题:基于Spark DataFrame的rec列生成分组ID

需求说明

现有Spark DataFrame,需基于rec列生成名为groupId的新列,规则为:每次遇到值为D的行时开启新分组。

示例输入

rec    amount  date           
D        250     20220522                  
C        110     20220522                  
D        120    20220522                   
C        100    20220522                   
C        50     20220522                   
D        50     20220522                   
D        50     20220522                   
D        50     20220522                   

期望输出

rec    amount  date            groupId   
D        250     20220522       1           
C        110     20220522       1           
D        120    20220522        2           
C        100    20220522        2          
C        50     20220522        2           
D        50     20220522        3           
D        50     20220522        4           
D        50     20220522        5           

尝试的错误代码

WindowSpec window = Window.orderBy("date");
Dataset<Row> dataset4 = data
            .withColumn("nextRow", functions.lead("rec", 1).over(window))
            .withColumn("prevRow", functions.lag("rec", 1).over(window))
            .withColumn("groupId",
                functions.when(functions.col("nextRow")
                        .equalTo(functions.col("prevRow")),
                        functions.dense_rank().over(window)
                    ));

错误分析

  1. 窗口排序逻辑错误:仅按date排序,但示例中所有行date值相同,Spark无法保证行的原始顺序,会导致分组逻辑混乱。
  2. 分组判断逻辑完全偏离需求:用nextRow和prevRow相等作为判断条件,和“遇到D开启新分组”的规则毫无关联,无法触发正确的分组切换。
  3. 分组ID计算错误:dense_rank()是基于排序的排名函数,不是累计计数新分组的次数,无法生成连续递增的分组ID。

正确实现方案

思路

  1. 生成行号保证原始行顺序(避免date相同时的排序混乱);
  2. 创建标记列:rec为D时标记为1,否则为0;
  3. 对标记列做累加求和,累加范围从第一行到当前行,求和结果即为groupId。

Java代码实现

import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.WindowSpec;
import static org.apache.spark.sql.functions.*;

// 1. 添加行号保证原始顺序
WindowSpec rowNumWindow = Window.orderBy(monotonically_increasing_id());
Dataset<Row> withRowNum = data.withColumn("row_num", row_number().over(rowNumWindow));

// 2. 创建分组标记列:rec为D时标记1,否则0
Dataset<Row> withFlag = withRowNum.withColumn("flag", when(col("rec").equalTo("D"), 1).otherwise(0));

// 3. 累加标记列生成groupId
WindowSpec groupWindow = Window.orderBy("row_num").rowsBetween(Window.unboundedPreceding, Window.currentRow);
Dataset<Row> result = withFlag.withColumn("groupId", sum("flag").over(groupWindow))
    .drop("row_num", "flag");

// 查看结果
result.show();

代码说明

  • monotonically_increasing_id()生成唯一递增ID,配合row_number()保证行的原始顺序;
  • 标记列flag在遇到D时触发累加,实现“每次D开启新分组”的规则;
  • 累加窗口rowsBetween(Window.unboundedPreceding, Window.currentRow)确保从第一行到当前行的标记值求和,得到连续的分组ID。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 22:30:58