This is an automated email from the ASF dual-hosted git repository.
ethanfeng pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/celeborn.git
The following commit(s) were added to refs/heads/main by this push:
new c84d0995a [CELEBORN-1399] MR CelebornMapOutputCollector should check
exception after flush
c84d0995a is described below
commit c84d0995a90090540286ff713b2ae8778ef4335e
Author: zhihu <[email protected]>
AuthorDate: Wed Apr 24 15:32:24 2024 +0800
[CELEBORN-1399] MR CelebornMapOutputCollector should check exception after
flush
### What changes were proposed in this pull request?
`CelebornMapOutputCollector` should check exception after flush for MR.
### Why are the changes needed?
For small shuffle mr job, shuffle push maybe only happend one time when map
task finished deal all task records and do flush before close
MapOutputCollector
In this case, CelebornMapOutputCollector should checkException after flush,
and throw exceptions when flush has exception, if not, job status is wrong
### Does this PR introduce _any_ user-facing change?
no
### How was this patch tested?
Test use mr job
Closes #2477 from lifulong/CELEBORN-1399.
Authored-by: zhihu <[email protected]>
Signed-off-by: mingji <[email protected]>
---
.../hadoop/mapred/CelebornMapOutputCollector.java | 6 +-
.../hadoop/mapred/CelebornSortBasedPusher.java | 2 +-
.../apache/celeborn/tests/mr/WordCountTest.scala | 69 ++++++++++++++++++++++
3 files changed, 74 insertions(+), 3 deletions(-)
diff --git
a/client-mr/mr/src/main/java/org/apache/hadoop/mapred/CelebornMapOutputCollector.java
b/client-mr/mr/src/main/java/org/apache/hadoop/mapred/CelebornMapOutputCollector.java
index fcecf85fb..1c2f57809 100644
---
a/client-mr/mr/src/main/java/org/apache/hadoop/mapred/CelebornMapOutputCollector.java
+++
b/client-mr/mr/src/main/java/org/apache/hadoop/mapred/CelebornMapOutputCollector.java
@@ -117,16 +117,18 @@ public class CelebornMapOutputCollector<K extends Object,
V extends Object>
}
@Override
- public void close() {
+ public void close() throws IOException {
logger.info("Mapper collector close");
reporter.progress();
celebornSortBasedPusher.close();
+ celebornSortBasedPusher.checkException();
}
@Override
- public void flush() {
+ public void flush() throws IOException {
logger.info("Mapper collector flush");
celebornSortBasedPusher.flush();
+ celebornSortBasedPusher.checkException();
reporter.progress();
}
}
diff --git
a/client-mr/mr/src/main/java/org/apache/hadoop/mapred/CelebornSortBasedPusher.java
b/client-mr/mr/src/main/java/org/apache/hadoop/mapred/CelebornSortBasedPusher.java
index 4692a58cf..69868197e 100644
---
a/client-mr/mr/src/main/java/org/apache/hadoop/mapred/CelebornSortBasedPusher.java
+++
b/client-mr/mr/src/main/java/org/apache/hadoop/mapred/CelebornSortBasedPusher.java
@@ -318,7 +318,7 @@ public class CelebornSortBasedPusher<K, V> extends
OutputStream {
numMappers);
shuffleClient.mapperEnd(0, mapId, attempt, numMappers);
} catch (IOException e) {
- logger.error("Mapper end failed, data lost", e);
+ exception.compareAndSet(null, e);
}
partitionedKVs.clear();
serializedKV = null;
diff --git
a/tests/mr-it/src/test/scala/org/apache/celeborn/tests/mr/WordCountTest.scala
b/tests/mr-it/src/test/scala/org/apache/celeborn/tests/mr/WordCountTest.scala
index d450edbb2..be8726d51 100644
---
a/tests/mr-it/src/test/scala/org/apache/celeborn/tests/mr/WordCountTest.scala
+++
b/tests/mr-it/src/test/scala/org/apache/celeborn/tests/mr/WordCountTest.scala
@@ -180,4 +180,73 @@ class WordCountTest extends AnyFunSuite with Logging with
MiniClusterFeature
assert(outputFilePath.toFile.exists())
assert(Files.readAllLines(outputFilePath).contains("celeborn\t1"))
}
+
+ test("celeborn mr integration test - word count shuffle exception") {
+ val input = Utils.createTempDir(System.getProperty("java.io.tmpdir"),
"input")
+ Files.write(
+ Paths.get(input.getPath, "v1.txt"),
+ "hello world celeborn".getBytes(StandardCharsets.UTF_8))
+ Files.write(
+ Paths.get(input.getPath, "v2.txt"),
+ "hello world mapreduce".getBytes(StandardCharsets.UTF_8))
+
+ val output = Utils.createTempDir(System.getProperty("java.io.tmpdir"),
"output")
+ val mrOutputPath = new Path(output.getPath + File.separator + "mr_output")
+
+ var exitCode = false
+ val conf = new Configuration(yarnCluster.getConfig)
+ // YARN config
+ conf.set("yarn.app.mapreduce.am.job.recovery.enable", "false")
+ conf.set(
+ "yarn.app.mapreduce.am.command-opts",
+ "org.apache.celeborn.mapreduce.v2.app.MRAppMasterWithCeleborn")
+
+ // MapReduce config
+ conf.set("mapreduce.framework.name", "yarn")
+ conf.set("mapreduce.job.user.classpath.first", "true")
+
+ conf.set("mapreduce.job.reduce.slowstart.completedmaps", "1")
+ conf.set(
+ "mapreduce.celeborn.master.endpoints",
+ s"errorhost:${master.conf.get(CelebornConf.MASTER_PORT)}")
+ conf.set(
+ MRJobConfig.MAP_OUTPUT_COLLECTOR_CLASS_ATTR,
+ "org.apache.hadoop.mapred.CelebornMapOutputCollector")
+ conf.set(
+ "mapreduce.job.reduce.shuffle.consumer.plugin.class",
+ "org.apache.hadoop.mapreduce.task.reduce.CelebornShuffleConsumer")
+
+ val job = Job.getInstance(conf, "word count")
+ job.setJarByClass(classOf[WordCount])
+ job.setMapperClass(classOf[WordCount.TokenizerMapper])
+ job.setCombinerClass(classOf[WordCount.IntSumReducer])
+ job.setReducerClass(classOf[WordCount.IntSumReducer])
+ job.setOutputKeyClass(classOf[Text])
+ job.setOutputValueClass(classOf[IntWritable])
+ FileInputFormat.addInputPath(job, new Path(input.getPath))
+ FileOutputFormat.setOutputPath(job, mrOutputPath)
+
+ val mapreduceLibPath =
+ (Utils.getCodeSourceLocation(getClass).split("/").dropRight(1) ++ Array(
+ "mapreduce_lib")).mkString("/")
+ val excludeJarList =
+ Seq(
+ "hadoop-client-api",
+ "hadoop-client-runtime",
+ "hadoop-client-minicluster",
+ "celeborn-client-mr-shaded",
+ "log4j")
+ Files.list(Paths.get(mapreduceLibPath)).iterator().asScala.foreach(path =>
{
+ if (!excludeJarList.exists(path.toFile.getPath.contains(_))) {
+ job.addFileToClassPath(new Path(path.toString))
+ }
+ })
+ logInfo(s"Job class path
${job.getFileClassPaths.map(_.toString).mkString(",")}")
+
+ exitCode = job.waitForCompletion(true)
+ assert(!exitCode, "Should return error code.")
+
+ val outputFilePath = Paths.get(mrOutputPath.toString, "part-r-00000")
+ assert(!outputFilePath.toFile.exists())
+ }
}