Spark 1.6中JavaRDD迭代出现无限循环问题求解决方案
问题分析与修复方案
你的代码之所以陷入无限循环,核心问题出在Spark的懒加载机制+循环依赖,再加上逻辑上的冗余,咱们一步步拆解:
为什么会无限循环?
你在map操作的闭包里调用了getMin(input)——这个方法肯定是要对input RDD执行action操作(比如reduce或者collect)来计算最小值对吧?而Spark的转换操作(比如map)是懒加载的,只有遇到action才会触发计算。现在的情况是:
- 当Spark要计算
map后的RDD时,每个元素的处理都需要先调用getMin(input),这会触发input的计算; - 但
input的计算又会触发它的转换链(也就是你的map操作),这样就形成了循环依赖:计算map需要input的min,计算input的min又需要先完成map,Spark就会一直卡在这个循环里,永远无法完成计算。
另外你的代码逻辑也有冗余:map把每个元素都转换成了min,之后filter判断的都是这些重复的min值,这显然和你“查找并删除被支配元素”的需求不符。
修复步骤
解决思路很简单:提前把min值计算出来,拿到Driver端作为本地变量,再用这个变量去处理RDD,避免在转换操作里重复触发RDD计算。
1. 先计算min值(用action操作拿到本地)
假设你的input是数值类型的JavaRDD(比如JavaRDD<Integer>),可以用Spark自带的reduce或者min方法直接计算:
// 计算input的最小值,这是一个action操作,会触发一次RDD计算并把结果拿到Driver端 Integer min = input.reduce((a, b) -> Math.min(a, b)); // 如果是其他可比类型,比如自定义对象,你可以调整reduce的逻辑,或者实现Comparator后用input.min(comparator)
2. 用本地的min值过滤元素
现在直接对原input RDD进行过滤,判断每个元素是否不被min支配,不需要多余的map操作:
JavaRDD<Integer> items = input.filter(x -> !Domine(x, min));
额外注意事项
- 永远不要在Spark的转换操作(
map/filter等)闭包里引用另一个RDD,或者依赖原RDD的action计算——这会导致重复计算甚至循环依赖; - 所有需要在转换中使用的聚合值,都要提前通过action操作拿到Driver端,作为本地变量传入闭包;
- 如果你的
Domine方法涉及复杂逻辑,确保它是可序列化的(因为Spark会把闭包里的变量序列化后发送到Executor端)。
内容的提问来源于stack exchange,提问作者user9467051
相关产品推荐
相关产品推荐

