PySpark中RDD的Resiliency特性与不可变性相关疑问
RDD不可变性相关疑问解答
你理解的偏差核心是混淆了变量引用和RDD实例本身的修改,你的代码完全符合RDD不可变性的要求,没有违反规则。
核心原理解释
- RDD的不可变性指的是:RDD实例对象在创建后,其内部存储的分区数据、依赖关系、计算逻辑等属性就完全固定,无法被任何操作修改。
- 你代码中用到的
map、reduceByKey、sortBy都属于RDD的转换算子,这类算子的逻辑从来不是修改原RDD,而是基于原RDD的计算逻辑生成一个全新的RDD实例作为返回值。 - 你反复对
filterData赋值的操作,本质是把filterData这个变量指针,一次次绑定到新生成的RDD实例上,原来的旧RDD实例如果没有其他变量引用,会被Spark的资源回收机制自动清理,全程没有任何RDD实例被修改。
我们可以把你的代码的本质拆解成以下逻辑,和你写的代码效果完全一致:
# 生成读取文件的初始RDD file_rdd = sc.textFile('file name') # map算子返回新的分割后RDD split_rdd = file_rdd.map(lambda x: x.split(',')) # reduceByKey算子返回新的聚合后RDD agg_rdd = split_rdd.reduceByKey(lambda x,y: x+y) # sortBy算子返回新的排序后RDD sorted_rdd = agg_rdd.sortBy(lambda x: x[1]) # 触发作业拉取最终结果 result = sorted_rdd.collect()
- 你提到的Resiliency(弹性)特性,正是基于RDD的不可变性实现的:Spark会记录每个RDD的血缘依赖关系,一旦作业运行中出现分区数据丢失,不需要全量重算所有数据,只需要根据血缘关系回溯计算对应丢失的分区即可。
和pandas DataFrame的差异补充
pandas的DataFrame虽然也有大量返回新对象的操作,但pandas支持通过inplace=True参数直接修改原DataFrame实例本身。而RDD的不可变性是强制约束,没有任何可以修改原RDD实例的操作,所有转换必然生成新的RDD。
内容的提问来源于stack exchange,提问作者fellowCoder
相关产品推荐
相关产品推荐

