Spark中基于Window函数分组,用最高rank值替换对应列内容
解决思路与实现代码
需求分析
按ID分组,找到每个ID对应的最高dateRank值,提取该行的PState和MState,并将该ID下所有行的这两个字段替换为提取的值。
实现步骤
- 计算每个
ID的最大dateRank值(需将字符串类型的dateRank转为整数类型后计算最大值); - 在每个
ID分区内,筛选出dateRank等于最大rank的行,提取对应的PState和MState; - 将提取到的目标值覆盖该
ID下所有行的原PState和MState字段,最后清理临时列。
Java代码实现
import org.apache.spark.sql.Dataset; import org.apache.spark.sql.Row; import org.apache.spark.sql.functions; import org.apache.spark.sql.types.DataTypes; import org.apache.spark.sql.Window; import org.apache.spark.sql.WindowSpec; // 1. 定义按ID分区的窗口,用于计算每个ID的最大rank WindowSpec idWindow = Window.partitionBy(functions.col("ID")); // 2. 添加max_rank列,计算每个ID的最高dateRank(转成Integer类型计算) Dataset<Row> withMaxRankDS = dateRankedDS.withColumn( "max_rank", functions.max(functions.col("dateRank").cast(DataTypes.IntegerType)).over(idWindow) ); // 3. 获取对应max_rank的PState和MState,替换原字段 Dataset<Row> resultDS = withMaxRankDS .withColumn( "target_PState", functions.first( functions.when( functions.col("dateRank").cast(DataTypes.IntegerType).equalTo(functions.col("max_rank")), functions.col("PState") ) ).over(idWindow) ) .withColumn( "target_MState", functions.first( functions.when( functions.col("dateRank").cast(DataTypes.IntegerType).equalTo(functions.col("max_rank")), functions.col("MState") ) ).over(idWindow) ) // 替换原字段 .withColumn("PState", functions.col("target_PState")) .withColumn("MState", functions.col("target_MState")) // 清理临时列 .drop("dateRank", "max_rank", "target_PState", "target_MState"); // 展示结果 resultDS.show(false);
输出结果
+---+------+-------+--------+------------------------+ |ID |PState|MState |dt |TS | +---+------+-------+--------+------------------------+ |1 |Iowa |CA |20240212|null | |1 |Iowa |CA |20250212|null | |2 |NJ |NY |20240212|2024-01-01T00:00:00.000Z| |2 |NJ |NY |20250212|null | |3 |CA |MS |20240212|null | |3 |CA |MS |20240212|null | |3 |CA |MS |20240212|null | +---+------+-------+--------+------------------------+
注:如果ID下有多个行对应最高rank(如ID3的所有行rank都是1),上述代码会取该ID分区内第一个符合条件的行的PState和MState。若需要取最后一个,可将
first替换为last;若需要自定义规则,可调整when内的逻辑。
内容的提问来源于stack exchange,提问作者Tim-Timer
相关产品推荐
相关产品推荐

