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

