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) ));
错误分析
- 窗口排序逻辑错误:仅按
date排序,但示例中所有行date值相同,Spark无法保证行的原始顺序,会导致分组逻辑混乱。 - 分组判断逻辑完全偏离需求:用
nextRow和prevRow相等作为判断条件,和“遇到D开启新分组”的规则毫无关联,无法触发正确的分组切换。 - 分组ID计算错误:
dense_rank()是基于排序的排名函数,不是累计计数新分组的次数,无法生成连续递增的分组ID。
正确实现方案
思路
- 生成行号保证原始行顺序(避免
date相同时的排序混乱); - 创建标记列:
rec为D时标记为1,否则为0; - 对标记列做累加求和,累加范围从第一行到当前行,求和结果即为
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
相关产品推荐
相关产品推荐

