为何Spark的filter操作无法保留分区机制?
你这个疑问特别合理——从直观逻辑看,filter就是在每个分区里挑符合条件的数据,分区数量好像也没变,为啥说它没法保留分区呢?其实这里的“保留分区”指的是无法保留原RDD的分区器(Partitioner)信息,而不是完全丢失分区结构,具体可以拆成这几点来看:
分区器的核心作用:Spark里的分区器(比如
HashPartitioner、RangePartitioner)是给键值对RDD(PairRDD)用的,它决定了数据的键值对该分配到哪个分区。当你执行filter时,虽然是在每个分区内筛选数据,但Spark没法保证筛选后的键值对还符合原来的分区规则。比如原RDD用HashPartitioner按key哈希分区,filter后某个key的所有数据都被过滤掉了,原分区器的规则对新RDD来说已经失去意义,所以Spark会主动丢弃原分区器。非PairRDD的特殊情况:如果是普通的非键值对RDD,本身就没有分区器,这里说的“无法保留分区”更多是指无法继承原分区的优化特性。比如原RDD是从HDFS读取的,分区对应HDFS块,filter后虽然分区数量不变,但每个分区的数据量可能差异极大,后续操作的数据本地性优化可能会打折扣,但分区本身的物理结构还是存在的。
对比能保留分区器的操作:比如
mapValues,这个操作只修改value而不碰key,key的哈希值没变,数据的分区规则完全有效,所以Spark可以放心保留原分区器。但filter可能改变每个分区内的键集合,甚至把某个分区的数据全滤掉,自然没法保留原分区器。
补充原文表述:部分操作(如map、flatMap、filter)无法保留分区。map、flatMap、filter操作会对每个分区执行函数。
这里的“无法保留分区”本质是指无法保留分区的元数据(尤其是分区器),而非分区的物理结构消失。
内容的提问来源于stack exchange,提问作者Hoori M.

