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

Spark中基于Window函数分组,用最高rank值替换对应列内容

解决思路与实现代码

需求分析

按ID分组,找到每个ID对应的最高dateRank值,提取该行的PState和MState,并将该ID下所有行的这两个字段替换为提取的值。

实现步骤

  1. 计算每个ID的最大dateRank值(需将字符串类型的dateRank转为整数类型后计算最大值);
  2. 在每个ID分区内,筛选出dateRank等于最大rank的行,提取对应的PState和MState;
  3. 将提取到的目标值覆盖该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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 04:52:04