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

