Spark中reduceByKey处理groupBy结果未达预期问题咨询
问题原因解析
你的问题核心在于**groupBy之后的RDD中每个key仅对应一个Iterable<Integer>值**,所以reduceByKey的合并逻辑根本不会执行。
具体拆解:
- 执行
groupBy(x -> x)后,Spark已经将所有相同key的元素聚合到同一个Iterable中,此时生成的JavaPairRDD<Integer, Iterable<Integer>>里,每个key在整个RDD中只有一条记录(不管原数据里该key有多少个元素)。 reduceByKey的作用是对同一个key下的多个值进行合并操作。但现在每个key只有一个值,没有其他值可以和它合并,所以Spark直接保留了原来的Iterable,完全不会调用你传入的Function2逻辑——不管你返回空列表、null还是原对象,都不会被执行。
验证方法:
你可以在合并函数里加个打印语句,比如:
(integers, integers2) -> { System.out.println("合并逻辑执行了!"); return new ArrayList<>(); }
运行后会发现控制台完全不会输出这句话,证明合并逻辑根本没被触发。
实现预期效果的正确方式:
如果想把每个key对应的Iterable替换成空列表,直接用mapValues即可,它会对每个key对应的value做转换:
JavaPairRDD<Integer, Iterable<Integer>> aa = integerIterableJavaPairRDD.mapValues(x -> new ArrayList<>()); aa.collect().forEach(System.out::println);
这样就能得到你预期的输出:
(1,[]) (2,[]) (3,[]) (4,[]) (5,[])
内容的提问来源于stack exchange,提问作者fei gao
相关产品推荐
相关产品推荐

