Author: ddas
Date: Sat May 14 01:15:53 2011
New Revision: 1102939

URL: http://svn.apache.org/viewvc?rev=1102939&view=rev
Log:
Support fail-fast for MR jobs. Contributed by Devaraj Das.

Modified:
    hadoop/mapreduce/branches/MR-279/CHANGES.txt
    
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/job/event/TaskAttemptEventType.java
    
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/job/impl/TaskAttemptImpl.java
    
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/job/impl/TaskImpl.java
    
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/rm/ContainerRequestEvent.java
    
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/rm/RMContainerAllocator.java

Modified: hadoop/mapreduce/branches/MR-279/CHANGES.txt
URL: 
http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/CHANGES.txt?rev=1102939&r1=1102938&r2=1102939&view=diff
==============================================================================
--- hadoop/mapreduce/branches/MR-279/CHANGES.txt (original)
+++ hadoop/mapreduce/branches/MR-279/CHANGES.txt Sat May 14 01:15:53 2011
@@ -3,6 +3,9 @@ Hadoop MapReduce Change Log
 Trunk (unreleased changes)
 
   MAPREDUCE-279
+   
+    Support fail-fast for MR jobs. (ddas)
+   
     Support mapreduce old (0.20) APIs. (sharad)
 
     Fixed reservation's bad interaction with delay scheduling in CS.

Modified: 
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/job/event/TaskAttemptEventType.java
URL: 
http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/job/event/TaskAttemptEventType.java?rev=1102939&r1=1102938&r2=1102939&view=diff
==============================================================================
--- 
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/job/event/TaskAttemptEventType.java
 (original)
+++ 
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/job/event/TaskAttemptEventType.java
 Sat May 14 01:15:53 2011
@@ -25,6 +25,7 @@ public enum TaskAttemptEventType {
 
   //Producer:Task
   TA_SCHEDULE,
+  TA_RESCHEDULE,
 
   //Producer:Client, Task
   TA_KILL,

Modified: 
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/job/impl/TaskAttemptImpl.java
URL: 
http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/job/impl/TaskAttemptImpl.java?rev=1102939&r1=1102938&r2=1102939&view=diff
==============================================================================
--- 
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/job/impl/TaskAttemptImpl.java
 (original)
+++ 
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/job/impl/TaskAttemptImpl.java
 Sat May 14 01:15:53 2011
@@ -158,7 +158,9 @@ public abstract class TaskAttemptImpl im
 
      // Transitions from the NEW state.
      .addTransition(TaskAttemptState.NEW, TaskAttemptState.UNASSIGNED,
-         TaskAttemptEventType.TA_SCHEDULE, new RequestContainerTransition())
+         TaskAttemptEventType.TA_SCHEDULE, new 
RequestContainerTransition(false))
+     .addTransition(TaskAttemptState.NEW, TaskAttemptState.UNASSIGNED,
+         TaskAttemptEventType.TA_RESCHEDULE, new 
RequestContainerTransition(true))
      .addTransition(TaskAttemptState.NEW, TaskAttemptState.KILLED,
          TaskAttemptEventType.TA_KILL, new KilledTransition())
      .addTransition(TaskAttemptState.NEW, TaskAttemptState.FAILED,
@@ -833,6 +835,10 @@ public abstract class TaskAttemptImpl im
   private static String[] racks = new String[] {NetworkTopology.DEFAULT_RACK};
   private static class RequestContainerTransition implements
       SingleArcTransition<TaskAttemptImpl, TaskAttemptEvent> {
+    boolean rescheduled = false;
+    public RequestContainerTransition(boolean rescheduled) {
+      this.rescheduled = rescheduled;
+    }
     @Override
     public void transition(TaskAttemptImpl taskAttempt, 
         TaskAttemptEvent event) {
@@ -840,10 +846,18 @@ public abstract class TaskAttemptImpl im
       taskAttempt.eventHandler.handle
           (new SpeculatorEvent(taskAttempt.getID().getTaskId(), +1));
       //request for container
-      taskAttempt.eventHandler.handle(
-          new ContainerRequestEvent(taskAttempt.attemptId, 
-              taskAttempt.resourceCapability, 
-              taskAttempt.getPriority(), taskAttempt.dataLocalHosts, racks));
+      if (rescheduled) {
+        taskAttempt.eventHandler.handle(
+            
ContainerRequestEvent.createContainerRequestEventForFailedContainer(
+                taskAttempt.attemptId, 
+                taskAttempt.resourceCapability, 
+                taskAttempt.getPriority()));
+      } else {
+        taskAttempt.eventHandler.handle(
+            new ContainerRequestEvent(taskAttempt.attemptId, 
+                taskAttempt.resourceCapability, 
+                taskAttempt.getPriority(), taskAttempt.dataLocalHosts, racks));
+      }
     }
   }
 

Modified: 
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/job/impl/TaskImpl.java
URL: 
http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/job/impl/TaskImpl.java?rev=1102939&r1=1102938&r2=1102939&view=diff
==============================================================================
--- 
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/job/impl/TaskImpl.java
 (original)
+++ 
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/job/impl/TaskImpl.java
 Sat May 14 01:15:53 2011
@@ -475,8 +475,13 @@ public abstract class TaskImpl implement
     ++nextAttemptNumber;
     ++numberUncompletedAttempts;
     //schedule the nextAttemptNumber
-    eventHandler.handle(new TaskAttemptEvent(attempt.getID(),
-        TaskAttemptEventType.TA_SCHEDULE));
+    if (failedAttempts > 0) {
+      eventHandler.handle(new TaskAttemptEvent(attempt.getID(),
+        TaskAttemptEventType.TA_RESCHEDULE));
+    } else {
+      eventHandler.handle(new TaskAttemptEvent(attempt.getID(),
+          TaskAttemptEventType.TA_SCHEDULE));
+    }
   }
 
   @Override

Modified: 
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/rm/ContainerRequestEvent.java
URL: 
http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/rm/ContainerRequestEvent.java?rev=1102939&r1=1102938&r2=1102939&view=diff
==============================================================================
--- 
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/rm/ContainerRequestEvent.java
 (original)
+++ 
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/rm/ContainerRequestEvent.java
 Sat May 14 01:15:53 2011
@@ -30,6 +30,7 @@ public class ContainerRequestEvent exten
   private Resource capability;
   private String[] hosts;
   private String[] racks;
+  private boolean earlierAttemptFailed = false;
 
   public ContainerRequestEvent(TaskAttemptId attemptID, 
       Resource capability, int priority,
@@ -41,6 +42,18 @@ public class ContainerRequestEvent exten
     this.hosts = hosts;
     this.racks = racks;
   }
+  
+  ContainerRequestEvent(TaskAttemptId attemptID, Resource capability, 
+      int priority) {
+    this(attemptID, capability, priority, null, null);
+    this.earlierAttemptFailed = true;
+  }
+  
+  public static ContainerRequestEvent 
createContainerRequestEventForFailedContainer(
+      TaskAttemptId attemptID, 
+      Resource capability, int priority) {
+    return new ContainerRequestEvent(attemptID,capability,priority);
+  }
 
   public Resource getCapability() {
     return capability;
@@ -57,4 +70,8 @@ public class ContainerRequestEvent exten
   public String[] getRacks() {
     return racks;
   }
+  
+  public boolean getEarlierAttemptFailed() {
+    return earlierAttemptFailed;
+  }
 }
\ No newline at end of file

Modified: 
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/rm/RMContainerAllocator.java
URL: 
http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/rm/RMContainerAllocator.java?rev=1102939&r1=1102938&r2=1102939&view=diff
==============================================================================
--- 
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/rm/RMContainerAllocator.java
 (original)
+++ 
hadoop/mapreduce/branches/MR-279/mr-client/hadoop-mapreduce-client-app/src/main/java/org/apache/hadoop/mapreduce/v2/app/rm/RMContainerAllocator.java
 Sat May 14 01:15:53 2011
@@ -129,20 +129,24 @@ public class RMContainerAllocator extend
     }
     eventList.add(event);
 
-    // Create resource requests
-    for (String host : event.getHosts()) {
-      // Data-local
-      addResourceRequest(event.getPriority(), host, event.getCapability());
-    }
+    if (event.getEarlierAttemptFailed()) {
+      addResourceRequest(event.getPriority(), ANY, event.getCapability());  
+    } else {
 
-    // Nothing Rack-local for now
-    for (String rack : event.getRacks()) {
-      addResourceRequest(event.getPriority(), rack, event.getCapability());
-    }
+      // Create resource requests
+      for (String host : event.getHosts()) {
+        // Data-local
+        addResourceRequest(event.getPriority(), host, event.getCapability());
+      }
 
-    // Off-switch
-    addResourceRequest(event.getPriority(), ANY, event.getCapability());
+      // Nothing Rack-local for now
+      for (String rack : event.getRacks()) {
+        addResourceRequest(event.getPriority(), rack, event.getCapability());
+      }
 
+      // Off-switch
+      addResourceRequest(event.getPriority(), ANY, event.getCapability());
+    }
   }
 
   private void addResourceRequest(Priority priority, String resourceName,
@@ -296,6 +300,13 @@ public class RMContainerAllocator extend
       Iterator<ContainerRequestEvent> it = requestList.iterator();
       while (it.hasNext()) {
         ContainerRequestEvent event = it.next();
+        if (event.getEarlierAttemptFailed()) {
+          // we want to fail fast. Ignore locality for rescheduling
+          // failed attempts.
+          assigned = event;
+          it.remove();
+          break;
+        }
         if (Arrays.asList(event.getHosts()).contains(host)) { // TODO: Fix
           assigned = event;
           it.remove();


Reply via email to