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

Spark Scala中小DataFrame与大DataFrame关联去重及特征选择问题

解决Spark Join后行数远超预期的问题

这个问题我在处理大表关联时碰过很多次,核心原因很明确:你的LargeDF和SmallDF里的cuid都存在重复记录。当执行inner join时,Small里每个重复的cuid都会和Large里所有同cuid的记录一一匹配,比如Small中某cuid出现3次,Large里同cuid有200条,join后就会生成3*200=600条数据,最终行数自然远超18万。

结合你提到的“无法按cuid分组,需要结合其他列做特征选择去重”的需求,给你几个可行的方案:

方案1:先对大表按「cuid+业务规则」去重,再关联

这是性能最优的方式,先把大表的重复数据砍掉,再和小表关联,能大幅减少后续的计算量。

核心思路是:针对每个cuid,根据业务逻辑选择唯一的一条记录(比如保留某列值最大/最新的行)。举个例子,如果你想保留每个cuid对应的Col3值最大的行:

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

// 对LargeDF按cuid分组,用row_number标记每个组内的行,保留Col3最大的那条
val deduplicatedLargeDF = largeDF
  .withColumn("row_rank", row_number().over(Window.partitionBy("cuid").orderBy(desc("Col3"))))
  .filter("row_rank = 1")
  .drop("row_rank")

// 再执行关联
val joinedDF = smallDF.join(deduplicatedLargeDF, Seq("cuid"), "inner")

你可以把orderBy(desc("Col3"))换成符合你业务需求的规则,比如按Col2的优先级排序、或者如果有时间列的话按最新时间排序——关键是要明确你想为每个cuid保留哪一条数据。

方案2:先清理小表的重复cuid(如果冗余的话)

如果SmallDF里的重复cuid是数据冗余(比如不需要保留重复的cuid条目),那可以先给小表去重,再关联:

// 去掉SmallDF中重复的cuid,每个cuid只留一条
val deduplicatedSmallDF = smallDF.dropDuplicates("cuid")

// 关联后的行数会等于LargeDF中匹配到的cuid行数(如果LargeDF还重复,行数还是会超过18万,这时候要结合方案1)
val joinedDF = deduplicatedSmallDF.join(largeDF, Seq("cuid"), "inner")

方案3:关联后再根据特征去重

如果必须先看到所有关联结果,再筛选去重,也可以在join后处理,但要注意这种方式的性能会差一些(因为要先处理大量的中间数据):

import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions.row_number

// 先执行关联
val joinedDF = smallDF.join(largeDF, Seq("cuid"), "inner")

// 对关联结果按cuid分组,选择每个组内的第一条(或按特征排序后的目标行)
val finalDF = joinedDF
  .withColumn("row_rank", row_number().over(Window.partitionBy("cuid").orderBy(desc("Col3"))))
  .filter("row_rank = 1")
  .drop("row_rank")

关于特征选择的关键提示

你提到要“结合其他列进行特征选择并去除重复”,本质是要明确什么是你认为的“重复数据”:

  • 如果是同一个cuid下,Col2+Col3的组合完全重复,那可以直接用dropDuplicates(Seq("cuid", "Col2", "Col3"))来去重;
  • 如果是同一个cuid下存在不同的列值,但你只需要保留一条,那必须用窗口函数(row_number/rank/dense_rank)结合业务规则来筛选唯一行。

另外,处理1900万行的大表时,建议调整一下分区数来优化性能,避免OOM:

// 根据你的集群资源设置合适的分区数,比如200-500之间
val optimizedLargeDF = largeDF.repartition(300)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.14 09:00:19