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

Databricks环境下基于同组最高频前驱值更新DataFrame/Hive表方案咨询

以下是适配Databricks环境的三种实现方案,均为分布式执行逻辑,可高效处理百万级及以上体量的数据:

1 PySpark实现

核心逻辑:先统计每个Name分组下shape的出现频次,取每个分组频次最高的shape值作为标准值,关联回原表完成替换。

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, count, row_number
from pyspark.sql.window import Window

# 1. 先计算每个分组下各shape值的频次
freq_df = df.groupBy("Name", "shape") \
            .agg(count("*").alias("freq"))

# 2. 取每个Name分组下频次最高的shape作为标准值
win = Window.partitionBy("Name").orderBy(col("freq").desc())
std_shape_df = freq_df.withColumn("rn", row_number().over(win)) \
                      .filter(col("rn") == 1) \
                      .select("Name", col("shape").alias("std_shape"))

# 3. 关联回原表,替换异常的shape值
result_df = df.join(std_shape_df, on="Name", how="left") \
              .drop("shape") \
              .withColumnRenamed("std_shape", "shape") \
              .select("sno", "Object", "Name", "shape", "rating")

如果需要替换多列,可按相同逻辑针对每列生成标准值后关联即可。

2 Hive SQL实现

直接通过窗口函数取分组高频值,关联完成替换:

-- 先计算每个分组的标准shape值
WITH shape_freq AS (
    SELECT 
        Name,
        shape,
        COUNT(*) AS freq,
        ROW_NUMBER() OVER(PARTITION BY Name ORDER BY COUNT(*) DESC) AS rn
    FROM original_table
    GROUP BY Name, shape
),
std_shape AS (
    SELECT Name, shape AS std_shape
    FROM shape_freq
    WHERE rn = 1
)
-- 关联替换
SELECT 
    t1.sno,
    t1.Object,
    t1.Name,
    t2.std_shape AS shape,
    t1.rating
FROM original_table t1
LEFT JOIN std_shape t2 ON t1.Name = t2.Name
3 Scala实现

逻辑和PySpark一致,语法适配Scala:

import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions.{col, count, row_number}

// 1. 统计分组频次
val freqDf = df.groupBy("Name", "shape")
               .agg(count("*").alias("freq"))

// 2. 取分组最高频次标准值
val win = Window.partitionBy("Name").orderBy(col("freq").desc)
val stdShapeDf = freqDf.withColumn("rn", row_number.over(win))
                       .filter(col("rn") === 1)
                       .select("Name", col("shape").alias("std_shape"))

// 3. 关联替换生成结果
val resultDf = df.join(stdShapeDf, Seq("Name"), "left")
                 .drop("shape")
                 .withColumnRenamed("std_shape", "shape")
                 .select("sno", "Object", "Name", "shape", "rating")

注意:如果同一分组下有多个值出现频次相同且都是最高,上述方案会随机取其中一个,若需要固定优先级可在窗口排序规则中额外补充排序字段。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 09:09:10