PySpark Fold算子理解困惑:结果不符及文档规则疑问
我正在理解PySpark中的Fold算子,实际输出和预期不符,也对文档描述有疑问。
先明确我对Fold的现有理解:
- Fold的结果依赖分区数量,因为
zeroValue会被应用(分区数+1)次:每个分区内应用一次,归约阶段从zeroValue开始合并分区结果。 - 默认分区数等于逻辑处理器数,我用的5600x是6核12线程,所以默认分区数12。
我用数值例子验证过:
sc.parallelize([1, 2, 3, 4, 5], 5).fold(3, lambda x,y:x+y)
执行流程:
每个分区1个元素,分别计算:
=> op(3, 1) = 4
=> op(3, 2) = 5
=> op(3, 3) = 6
=> op(3, 4) = 7
=> op(3, 5) = 8
归约阶段从zeroValue=3开始合并:
3 op 4 =7 →7 op5=12 →12 op6=18 →18 op7=25 →25 op8=33,最终结果33,符合预期。
但测试文档描述时出现问题,文档原文:
The function
op(t1, t2)is allowed to modifyt1and return it as its result value to avoid object allocation; however, it should not modifyt2.
我写了如下代码:
class Fold_Class: def __init__(self, alpha=0): self.alpha = alpha def fold_function(obj1, obj2): obj2.alpha *= 2 return obj2 ls = sc.parallelize([Fold_Class(1), Fold_Class(2), Fold_Class(3)], 3) folded_obj = ls.fold(Fold_Class(0), fold_function) print(folded_obj.alpha)
输出是12,但我预期是24,我的推导:
=> op(Fold_Class(0), Fold_Class(1)) → Fold_Class(2)
=> op(Fold_Class(0), Fold_Class(2)) → Fold_Class(4)
=> op(Fold_Class(0), Fold_Class(3)) → Fold_Class(6)
归约:
=> reduce(Fold_Class(2), Fold_Class(4)) → Fold_Class(8)
=> reduce(Fold_Class(8), Fold_Class(6)) → Fold_Class(12)
最后op(Fold_Class(0), Fold_Class(12)) → Fold_Class(24)
现在有两个问题:
- 我的
fold_function修改了obj2,根据文档这是不允许的,我是否理解错了文档内容? - 为什么fold函数返回的结果是12?
问题1解答
你没有理解错文档内容:你的fold_function确实违反了文档的要求。
文档明确说明op(t1,t2)不应该修改t2,原因在于Spark的RDD是不可变(immutable)的分布式数据集,每个分区的元素是被设计为只读的。如果修改t2,会破坏RDD的不可变性契约:比如当RDD被缓存、复用或者重新计算时,被修改过的t2会导致后续计算出现不可预测的结果,引发难以排查的bug。
允许修改t1是因为t1通常是归约过程中的累加器(由zeroValue衍生的实例),属于计算过程中的临时状态,修改它可以避免频繁创建新对象,提升性能,但t2作为RDD中的原始元素,绝对不能被修改。
问题2解答
你的推导错误在于对Fold归约阶段的流程理解偏差,正确的执行流程如下:
Spark的Fold执行分为两步:
- 分区内折叠:每个分区使用克隆后的
zeroValue作为初始值,对分区内元素执行foldLeft(即从左到右依次应用op)。 - 全局归约:从克隆后的
zeroValue开始,依次将每个分区的折叠结果传入op进行合并,而不是先合并所有分区结果再和zeroValue做一次op。
对应你的代码,具体执行步骤:
分区内折叠(3个分区,每个分区1个元素)
- 分区1:
foldLeft(Fold_Class(0))→ 执行op(Fold_Class(0), Fold_Class(1)),修改obj2的alpha为2,返回该实例 → 分区结果:Fold_Class(2) - 分区2:
foldLeft(Fold_Class(0))→ 执行op(Fold_Class(0), Fold_Class(2)),修改obj2的alpha为4,返回该实例 → 分区结果:Fold_Class(4) - 分区3:
foldLeft(Fold_Class(0))→ 执行op(Fold_Class(0), Fold_Class(3)),修改obj2的alpha为6,返回该实例 → 分区结果:Fold_Class(6)
全局归约(从Fold_Class(0)开始合并分区结果)
- 第一步:
op(Fold_Class(0), 分区1结果)→ 修改分区1结果的alpha为2*2=4,返回该实例 → 当前全局结果:Fold_Class(4) - 第二步:
op(Fold_Class(4), 分区2结果)→ 修改分区2结果的alpha为4*2=8,返回该实例 → 当前全局结果:Fold_Class(8) - 第三步:
op(Fold_Class(8), 分区3结果)→ 修改分区3结果的alpha为6*2=12,返回该实例 → 最终全局结果:Fold_Class(12)
这就是最终输出为12的原因。
内容的提问来源于stack exchange,提问作者Shashanka Das

