Spark作业mapPartitions重复执行问题及单次执行实现咨询
解决Spark作业mapPartitions重复执行的问题
问题根源
虽然对tuple2Dataset调用了cache(),但Spark的缓存是懒加载机制——只有当遇到行动算子(如count()、write())时才会实际执行计算并缓存数据。你的代码中先后执行了successfulRows.write()和failedRows.write(),这两个行动算子会分别触发对tuple2Dataset的计算,导致mapPartitions被执行两次。
修复方案
在缓存tuple2Dataset后,主动调用一个行动算子触发缓存落地,确保后续的读写操作直接使用缓存数据,而不是重新计算。常用的行动算子是count():
Dataset<Tuple2<Boolean, String>> tuple2Dataset = rows.mapPartitions(new GetPaymentsGraphFeatures(accumulator), Encoders.tuple(Encoders.BOOLEAN(), Encoders.STRING())); tuple2Dataset.cache(); // 主动触发缓存,执行一次mapPartitions并将结果存入缓存 tuple2Dataset.count();
修改后的完整核心代码片段:
Dataset<Tuple2<Boolean, String>> tuple2Dataset = rows.mapPartitions(new GetPaymentsGraphFeatures(accumulator), Encoders.tuple(Encoders.BOOLEAN(), Encoders.STRING())); tuple2Dataset.cache(); // 新增触发缓存的行动算子 tuple2Dataset.count(); Dataset<Row> successfulRows = tuple2Dataset.filter((FilterFunction<Tuple2<Boolean, String>>) booleanRowTuple2 -> booleanRowTuple2._1).map( (MapFunction<Tuple2<Boolean, String>, Row>) booleanRowTuple2 -> mapToRow(booleanRowTuple2._2), RowEncoder.apply(getSchema())); Dataset<Row> failedRows = tuple2Dataset.filter((FilterFunction<Tuple2<Boolean, String>>) booleanRowTuple2 -> !booleanRowTuple2._1).map( (MapFunction<Tuple2<Boolean, String>, Row>) booleanRowTuple2 -> mapToRow(booleanRowTuple2._2), RowEncoder.apply(getFailureSchema())); successfulRows.write().mode("overwrite").saveAsTable("db.deepak_jan_result"); failedRows.write().mode("overwrite").saveAsTable("db.deepak_jan_result_failures"); tuple2Dataset.unpersist();
额外优化建议
- 若集群内存资源有限,可使用
persist(StorageLevel.MEMORY_AND_DISK())替代cache(),当内存不足时数据会写入磁盘,避免缓存溢出导致重新计算。 - 检查
GetPaymentsGraphFeatures中的逻辑,确保除了Spark安全的累加器外,没有依赖其他外部可变状态,防止缓存数据出现不一致。
内容的提问来源于stack exchange,提问作者Deepak Patankar
相关产品推荐
相关产品推荐

