Spark Dataset使用自定义类执行reduce操作返回结果异常如何解决
问题根因
- 异常来自
Dataset构造阶段的map逻辑错误:你复用了同一个MapReducer实例sampleMR处理所有输入Row,没有为每条输入生成独立的结果对象。 - 你的
incrementCount方法逻辑是修改当前实例的count属性后返回this,Spark会在每次map调用后用Encoder序列化返回结果再处理下一条数据,这就导致第N条Row对应的输出MapReducer的count值为N,而非预期的1。比如某分区有5条数据,该分区输出的5个元素count值分别为1、2、3、4、5,求和后比正确值5多10,你遇到的correct_value +5是该异常逻辑在你的实际分区分布下的计算结果。 - 你定义的reduce相加逻辑、
MapReducer类的序列化和Bean规范都是符合要求的,不需要修改。
修复方法
修复map阶段逻辑
你只需要保证每条输入Row对应一个count=1的独立MapReducer对象即可,有两种修改方案:
方案1:每次处理Row时新建独立实例(推荐,逻辑更清晰)
Dataset<MapReducer> df = InitialTable.quartetsTable .toDF() .map( (MapFunction<Row, MapReducer>) row -> { MapReducer mr = new MapReducer(); mr.incrementCount(row); return mr; }, Encoders.bean(MapReducer.class) );
方案2:修改incrementCount方法,返回新实例而非修改自身
// 修改MapReducer类中的incrementCount实现 public MapReducer incrementCount(Row row){ MapReducer newMr = new MapReducer(); newMr.count = this.count + 1; return newMr; }
修改完成后保持原有reduce逻辑不变,即可得到正确的计数结果。
分区行数统计优化(可选)
如果你的需求是统计各分区的行数而非总行数,直接使用mapPartitions算子效率更高,无需走reduce流程:
Dataset<Integer> partitionCounts = InitialTable.quartetsTable .toDF() .mapPartitions(iterator -> { int cnt = 0; while (iterator.hasNext()) { iterator.next(); cnt++; } return Collections.singletonList(cnt).iterator(); }, Encoders.INT()); // 采集后得到所有分区的行数列表 List<Integer> cntList = partitionCounts.collectAsList();
内容的提问来源于stack exchange,提问作者Shariar Kabir
相关产品推荐
相关产品推荐

