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();