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

Apache Spark中Dataset.dropDuplicates()是否保留分区?

Spark 2.2 Dataset中dropDuplicates是否保留父数据集分区?

好问题!针对你提到的场景——已经按one、two分区并缓存的Dataset调用dropDuplicates(List("one", "two")),结论很明确:dropDuplicates不会保留父数据集的分区结构,哪怕你的去重键和分区键完全一致。

原因分析

Spark的dropDuplicates(底层对应distinct操作)的核心逻辑是要确保所有重复的记录被聚在一起,才能准确完成去重。这意味着无论父数据集的分区策略是什么,dropDuplicates都会触发shuffle操作——它会根据你指定的去重键重新分区数据,原来的分区信息会被完全覆盖。

在你的示例中,虽然你提前用repartition('one, 'two)把数据按one和two分了区,但Spark无法保证每个分区内已经包含了所有相同(one, two)的记录(除非你能确保数据本身是全局按这两个键分区且没有跨分区的重复),所以必须通过shuffle来全局聚合相同键的数据,才能完成准确去重。

代码验证

你可以在代码中添加分区查看逻辑,直观确认这一点:

case class Item(one: Int, two: Int, three: Int)
import session.implicits._

val ds = session.createDataset(List(Item(1,2,3), Item(1,2,3)))
val repart = ds.repartition('one, 'two).cache()

// 查看原分区信息
println(s"repart的分区数:${repart.rdd.getNumPartitions}")
repart.rdd.glom().foreach(arr => println(s"分区数据:${arr.mkString(",")}"))

// 执行去重并查看结果分区
val dedupDs = repart.dropDuplicates(List("one", "two"))
println(s"\n去重后的分区数:${dedupDs.rdd.getNumPartitions}")
dedupDs.rdd.glom().foreach(arr => println(s"分区数据:${arr.mkString(",")}"))

运行后你会发现,去重后的分区数、数据分布和原repart的分区完全不同,这就是shuffle带来的变化。

优化小技巧

如果想减少不必要的shuffle开销,可以先在每个分区内本地去重,再执行全局去重,这样能大幅减少shuffle的数据量:

// 先在分区内去重
val localDedup = repart.mapPartitions(iter => iter.distinct)
// 再全局去重
val finalDedup = localDedup.dropDuplicates(List("one", "two"))

不过要注意,最后的全局去重仍然会触发shuffle,只是数据量更小了。

内容的提问来源于stack exchange,提问作者Peter Jakubčo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:22:50