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
相关产品推荐
相关产品推荐

