This is an automated email from the ASF dual-hosted git repository.

nicholasjiang pushed a commit to branch branch-0.4
in repository https://gitbox.apache.org/repos/asf/celeborn.git

commit d6a60ace327ae2766e12d61e79fae51f3fb687ae
Author: lifulong <[email protected]>
AuthorDate: Wed Apr 24 15:32:24 2024 +0800

    [CELEBORN-1399] MR CelebornMapOutputCollector should check exception after 
flush
    
    `CelebornMapOutputCollector` should check exception after flush for MR.
    
    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
    
    no
    
    Test use mr job
    
    Closes #2477 from lifulong/CELEBORN-1399.
    
    Authored-by: lifulong <[email protected]>
    Signed-off-by: mingji <[email protected]>
    (cherry picked from commit c84d0995a90090540286ff713b2ae8778ef4335e)
    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 4148733d9..e1b89e269 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
@@ -151,4 +151,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())
+  }
 }

Reply via email to