Modified: hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/fifo/FifoScheduler.java URL: http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/fifo/FifoScheduler.java?rev=1084866&r1=1084865&r2=1084866&view=diff ============================================================================== --- hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/fifo/FifoScheduler.java (original) +++ hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/fifo/FifoScheduler.java Thu Mar 24 07:52:35 2011 @@ -22,9 +22,11 @@ import java.io.IOException; import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.Arrays; +import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.TreeMap; + import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.apache.hadoop.classification.InterfaceAudience.LimitedPrivate; @@ -32,24 +34,24 @@ import org.apache.hadoop.classification. import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.net.Node; import org.apache.hadoop.security.UserGroupInformation; +import org.apache.hadoop.yarn.ApplicationID; +import org.apache.hadoop.yarn.Container; +import org.apache.hadoop.yarn.ContainerToken; +import org.apache.hadoop.yarn.NodeID; +import org.apache.hadoop.yarn.Priority; +import org.apache.hadoop.yarn.Resource; +import org.apache.hadoop.yarn.ResourceRequest; import org.apache.hadoop.yarn.security.ContainerTokenIdentifier; +import org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ASMEvent; +import org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ApplicationMasterEvents.ApplicationTrackerEventType; import org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.NodeInfo; import org.apache.hadoop.yarn.server.resourcemanager.scheduler.Application; -import org.apache.hadoop.yarn.server.resourcemanager.scheduler.ClusterTracker; -import org.apache.hadoop.yarn.server.resourcemanager.scheduler.ClusterTrackerImpl; import org.apache.hadoop.yarn.server.resourcemanager.scheduler.NodeManager; +import org.apache.hadoop.yarn.server.resourcemanager.scheduler.NodeResponse; import org.apache.hadoop.yarn.server.resourcemanager.scheduler.NodeType; import org.apache.hadoop.yarn.server.resourcemanager.scheduler.Queue; import org.apache.hadoop.yarn.server.resourcemanager.scheduler.ResourceScheduler; -import org.apache.hadoop.yarn.server.resourcemanager.scheduler.ClusterTracker.NodeResponse; import org.apache.hadoop.yarn.server.security.ContainerTokenSecretManager; -import org.apache.hadoop.yarn.ApplicationID; -import org.apache.hadoop.yarn.Container; -import org.apache.hadoop.yarn.ContainerToken; -import org.apache.hadoop.yarn.NodeID; -import org.apache.hadoop.yarn.Priority; -import org.apache.hadoop.yarn.Resource; -import org.apache.hadoop.yarn.ResourceRequest; @LimitedPrivate("yarn") @Evolving @@ -59,7 +61,6 @@ public class FifoScheduler implements Re Configuration conf; private ContainerTokenSecretManager containerTokenSecretManager; - private final ClusterTracker clusterTracker; // TODO: The memory-block size should be site-configurable? public static final int MINIMUM_MEMORY = 1024; @@ -81,20 +82,14 @@ public class FifoScheduler implements Re } }; - public FifoScheduler() { - this.clusterTracker = createClusterTracker(); - } + public FifoScheduler() {} public FifoScheduler(Configuration conf, ContainerTokenSecretManager containerTokenSecretManager) { - this(); reinitialize(conf, containerTokenSecretManager); } - protected ClusterTracker createClusterTracker() { - return new ClusterTrackerImpl(); - } @Override public void reinitialize(Configuration conf, @@ -141,7 +136,7 @@ public class FifoScheduler implements Re private void releaseContainers(Application application, List<Container> release) { application.releaseContainers(release); for (Container container : release) { - clusterTracker.releaseContainer(application.getApplicationId(), container); + releaseContainer(application.getApplicationId(), container); } } @@ -161,7 +156,6 @@ public class FifoScheduler implements Re return applications.get(applicationId); } - @Override public synchronized void addApplication(ApplicationID applicationId, String user, String unusedQueue, Priority unusedPriority) throws IOException { @@ -171,7 +165,6 @@ public class FifoScheduler implements Re ", currently active: " + applications.size()); } - @Override public synchronized void removeApplication(ApplicationID applicationId) throws IOException { Application application = getApplication(applicationId); @@ -184,7 +177,7 @@ public class FifoScheduler implements Re releaseContainers(application, application.getCurrentContainers()); // Let the cluster know that the applications are done - clusterTracker.finishedApplication(applicationId, + finishedApplication(applicationId, application.getAllNodesForApplication()); // Remove the application @@ -389,9 +382,8 @@ public class FifoScheduler implements Re containers.add(container); } application.allocate(type, node, priority, request, containers); - clusterTracker.addAllocatedContainers(node, application.getApplicationId(), containers); + addAllocatedContainers(node, application.getApplicationId(), containers); } - return assignedContainers; } @@ -412,7 +404,7 @@ public class FifoScheduler implements Re public synchronized NodeResponse nodeUpdate(NodeInfo node, Map<CharSequence,List<Container>> containers ) { - NodeResponse nodeResponse = clusterTracker.nodeUpdate(node, containers); + NodeResponse nodeResponse = nodeUpdateInternal(node, containers); applicationCompletedContainers(nodeResponse.getCompletedContainers()); LOG.info("Node heartbeat " + node.getNodeID() + " resource = " + node.getAvailableResource()); if (org.apache.hadoop.yarn.server.resourcemanager.resource.Resource. @@ -426,15 +418,88 @@ public class FifoScheduler implements Re // preemption. return nodeResponse; } - + @Override - public NodeInfo addNode(NodeID nodeId,String hostName, - Node node, Resource capability) { - return clusterTracker.addNode(nodeId, hostName, node, capability); + public synchronized void handle(ASMEvent<ApplicationTrackerEventType> event) { + switch(event.getType()) { + case ADD: + try { + addApplication(event.getAppContext().getApplicationID(), event.getAppContext().getUser(), + event.getAppContext().getQueue(), event.getAppContext().getSubmissionContext().priority); + } catch(IOException ie) { + LOG.error("Unable to add application " + event.getAppContext().getApplicationID(), ie); + /** this is fatal we are not able to add applications for scheduling **/ + //TODO handle it later. + } + break; + case REMOVE: + try { + + removeApplication(event.getAppContext().getApplicationID()); + } catch(IOException ie) { + LOG.error("Unable to remove application " + event.getAppContext().getApplicationID(), ie); + } + break; + } + } + + private Map<String, NodeManager> nodes = new HashMap<String, NodeManager>(); + private Resource clusterResource = new Resource(); + + public synchronized Resource getClusterResource() { + return clusterResource; } @Override - public void removeNode(NodeInfo node) { - clusterTracker.removeNode(node); + public synchronized void removeNode(NodeInfo nodeInfo) { + org.apache.hadoop.yarn.server.resourcemanager.resource.Resource.subtractResource( + clusterResource, nodeInfo.getTotalCapability()); + nodes.remove(nodeInfo.getHostName()); + } + + public synchronized boolean isTracked(NodeInfo nodeInfo) { + NodeManager node = nodes.get(nodeInfo.getHostName()); + return (node == null? false: true); + } + + @Override + public synchronized NodeInfo addNode(NodeID nodeId, + String hostName, Node node, Resource capability) { + NodeManager nodeManager = new NodeManager(nodeId, hostName, node, capability); + nodes.put(nodeManager.getHostName(), nodeManager); + org.apache.hadoop.yarn.server.resourcemanager.resource.Resource.addResource( + clusterResource, nodeManager.getTotalCapability()); + return nodeManager; + } + + public synchronized boolean releaseContainer(ApplicationID applicationId, + Container container) { + // Reap containers + LOG.info("Application " + applicationId + " released container " + container); + NodeManager nodeManager = nodes.get(container.hostName.toString()); + return nodeManager.releaseContainer(container); + } + + private synchronized NodeResponse nodeUpdateInternal(NodeInfo nodeInfo, + Map<CharSequence,List<Container>> containers) { + NodeManager node = nodes.get(nodeInfo.getHostName()); + LOG.debug("nodeUpdate: node=" + nodeInfo.getHostName() + + " available=" + nodeInfo.getAvailableResource().memory); + return node.statusUpdate(containers); + + } + + public synchronized void addAllocatedContainers(NodeInfo nodeInfo, + ApplicationID applicationId, List<Container> containers) { + NodeManager node = nodes.get(nodeInfo.getHostName()); + node.allocateContainer(applicationId, containers); + } + + public synchronized void finishedApplication(ApplicationID applicationId, + List<NodeInfo> nodesToNotify) { + for (NodeInfo node: nodesToNotify) { + NodeManager nodeManager = nodes.get(node.getHostName()); + nodeManager.notifyFinishedApplication(applicationId); + } } }
Modified: hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMLaunchFailure.java URL: http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMLaunchFailure.java?rev=1084866&r1=1084865&r2=1084866&view=diff ============================================================================== --- hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMLaunchFailure.java (original) +++ hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMLaunchFailure.java Thu Mar 24 07:52:35 2011 @@ -1,26 +1,27 @@ /** -* Licensed to the Apache Software Foundation (ASF) under one -* or more contributor license agreements. See the NOTICE file -* distributed with this work for additional information -* regarding copyright ownership. The ASF licenses this file -* to you under the Apache License, Version 2.0 (the -* "License"); you may not use this file except in compliance -* with the License. You may obtain a copy of the License at -* -* http://www.apache.org/licenses/LICENSE-2.0 -* -* Unless required by applicable law or agreed to in writing, software -* distributed under the License is distributed on an "AS IS" BASIS, -* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -* See the License for the specific language governing permissions and -* limitations under the License. -*/ + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ package org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager; import java.io.IOException; import java.util.Arrays; import java.util.List; +import java.util.concurrent.atomic.AtomicInteger; import junit.framework.TestCase; @@ -63,19 +64,9 @@ public class TestAMLaunchFailure extends @Override public List<Container> allocate(ApplicationID applicationId, - List<ResourceRequest> ask, List<Container> release) throws IOException { + List<ResourceRequest> ask, List<Container> release) throws IOException { return Arrays.asList(container); } - - @Override - public void addApplication(ApplicationID applicationId, String user, - String queue, Priority priority) throws IOException { - } - - @Override - public void removeApplication(ApplicationID applicationId) - throws IOException { - } } private class DummyApplicationTracker implements EventHandler<ASMEvent<ApplicationTrackerEventType>> { @@ -90,9 +81,9 @@ public class TestAMLaunchFailure extends public class ExtApplicationsManagerImpl extends ApplicationsManagerImpl { private class DummyApplicationMasterLauncher implements EventHandler<ASMEvent<AMLauncherEventType>> { - private Object notify = new Object(); + private AtomicInteger notify = new AtomicInteger(); private AppContext app; - + public DummyApplicationMasterLauncher(ASMContext context) { context.getDispatcher().register(AMLauncherEventType.class, this); new TestThread().start(); @@ -101,8 +92,10 @@ public class TestAMLaunchFailure extends public void handle(ASMEvent<AMLauncherEventType> appEvent) { switch(appEvent.getType()) { case LAUNCH: + LOG.info("LAUNCH called "); app = appEvent.getAppContext(); synchronized (notify) { + notify.addAndGet(1); notify.notify(); } break; @@ -113,27 +106,29 @@ public class TestAMLaunchFailure extends public void run() { synchronized(notify) { try { - notify.wait(); + while (notify.get() == 0) { + notify.wait(); + } } catch (InterruptedException e) { e.printStackTrace(); } context.getDispatcher().getEventHandler().handle( - new ASMEvent<ApplicationEventType>(ApplicationEventType.LAUNCHED, - app)); + new ASMEvent<ApplicationEventType>(ApplicationEventType.LAUNCHED, + app)); } } } } public ExtApplicationsManagerImpl( - ApplicationTokenSecretManager applicationTokenSecretManager, - YarnScheduler scheduler) { + ApplicationTokenSecretManager applicationTokenSecretManager, + YarnScheduler scheduler) { super(applicationTokenSecretManager, scheduler, context); } @Override protected EventHandler<ASMEvent<AMLauncherEventType>> createNewApplicationMasterLauncher( - ApplicationTokenSecretManager tokenSecretManager) { + ApplicationTokenSecretManager tokenSecretManager) { return new DummyApplicationMasterLauncher(context); } } @@ -146,6 +141,7 @@ public class TestAMLaunchFailure extends Configuration conf = new Configuration(); new DummyApplicationTracker(); conf.setLong(YarnConfiguration.AM_EXPIRY_INTERVAL, 3000L); + conf.setInt(YarnConfiguration.AM_MAX_RETRIES, 1); asmImpl.init(conf); asmImpl.start(); } Modified: hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMRMRPCResponseId.java URL: http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMRMRPCResponseId.java?rev=1084866&r1=1084865&r2=1084866&view=diff ============================================================================== --- hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMRMRPCResponseId.java (original) +++ hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMRMRPCResponseId.java Thu Mar 24 07:52:35 2011 @@ -75,14 +75,6 @@ public class TestAMRMRPCResponseId exten List<ResourceRequest> ask, List<Container> release) throws IOException { return null; } - @Override - public void addApplication(ApplicationID applicationId, String user, - String queue, Priority priority) throws IOException { - } - @Override - public void removeApplication(ApplicationID applicationId) - throws IOException { - } } @Before Added: hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMRestart.java URL: http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMRestart.java?rev=1084866&view=auto ============================================================================== --- hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMRestart.java (added) +++ hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMRestart.java Thu Mar 24 07:52:35 2011 @@ -0,0 +1,300 @@ +package org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager; + +import java.io.IOException; +import java.nio.ByteBuffer; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.atomic.AtomicInteger; + +import junit.framework.TestCase; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.net.Node; +import org.apache.hadoop.yarn.ApplicationID; +import org.apache.hadoop.yarn.ApplicationMaster; +import org.apache.hadoop.yarn.ApplicationState; +import org.apache.hadoop.yarn.ApplicationStatus; +import org.apache.hadoop.yarn.ApplicationSubmissionContext; +import org.apache.hadoop.yarn.Container; +import org.apache.hadoop.yarn.ContainerID; +import org.apache.hadoop.yarn.ContainerToken; +import org.apache.hadoop.yarn.NodeID; +import org.apache.hadoop.yarn.Resource; +import org.apache.hadoop.yarn.ResourceRequest; +import org.apache.hadoop.yarn.conf.YarnConfiguration; +import org.apache.hadoop.yarn.event.EventHandler; +import org.apache.hadoop.yarn.security.ApplicationTokenSecretManager; +import org.apache.hadoop.yarn.server.resourcemanager.ResourceManager; +import org.apache.hadoop.yarn.server.resourcemanager.ResourceManager.ASMContext; +import org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ASMEvent; +import org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ApplicationMasterEvents.AMLauncherEventType; +import org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ApplicationMasterEvents.ApplicationEventType; +import org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ApplicationMasterEvents.ApplicationTrackerEventType; +import org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.NodeInfo; +import org.apache.hadoop.yarn.server.resourcemanager.scheduler.NodeResponse; +import org.apache.hadoop.yarn.server.resourcemanager.scheduler.ResourceScheduler; +import org.apache.hadoop.yarn.server.resourcemanager.scheduler.YarnScheduler; +import org.apache.hadoop.yarn.server.security.ContainerTokenSecretManager; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; + +/** + * Test to restart the AM on failure. + * + */ +public class TestAMRestart extends TestCase { + private static final Log LOG = LogFactory.getLog(TestAMRestart.class); + ApplicationsManagerImpl appImpl; + ASMContext asmContext = new ResourceManager.ASMContextImpl(); + ApplicationTokenSecretManager appTokenSecretManager = + new ApplicationTokenSecretManager(); + DummyResourceScheduler scheduler; + int count = 0; + ApplicationID appID; + final int maxFailures = 3; + AtomicInteger launchNotify = new AtomicInteger(); + AtomicInteger schedulerNotify = new AtomicInteger(); + volatile boolean stop = false; + int schedulerAddApplication = 0; + int schedulerRemoveApplication = 0; + int launcherLaunchCalled = 0; + int launcherCleanupCalled = 0; + ApplicationMasterInfo masterInfo; + + private class ExtApplicationsManagerImpl extends ApplicationsManagerImpl { + public ExtApplicationsManagerImpl( + ApplicationTokenSecretManager applicationTokenSecretManager, + YarnScheduler scheduler, ASMContext asmContext) { + super(applicationTokenSecretManager, scheduler, asmContext); + } + + @Override + public EventHandler<ASMEvent<AMLauncherEventType>> createNewApplicationMasterLauncher( + ApplicationTokenSecretManager tokenSecretManager) { + return new DummyAMLauncher(); + } + } + + private class DummyAMLauncher implements EventHandler<ASMEvent<AMLauncherEventType>> { + + public DummyAMLauncher() { + asmContext.getDispatcher().register(AMLauncherEventType.class, this); + new Thread() { + public void run() { + while (!stop) { + LOG.info("DEBUG -- waiting for launch"); + synchronized(launchNotify) { + while (launchNotify.get() == 0) { + try { + launchNotify.wait(); + } catch (InterruptedException e) { + } + } + asmContext.getDispatcher().getEventHandler().handle(new + ASMEvent<ApplicationEventType>(ApplicationEventType.LAUNCHED, + new TestAppContext(appID))); + launchNotify.addAndGet(-1); + } + } + } + }.start(); + } + + @Override + public void handle(ASMEvent<AMLauncherEventType> event) { + switch (event.getType()) { + case CLEANUP: + launcherCleanupCalled++; + break; + case LAUNCH: + LOG.info("DEBUG -- launching"); + launcherLaunchCalled++; + synchronized (launchNotify) { + launchNotify.addAndGet(1); + launchNotify.notify(); + } + break; + default: + break; + } + } + } + + private class DummyResourceScheduler implements ResourceScheduler { + @Override + public NodeInfo addNode(NodeID nodeId, String hostName, Node node, + Resource capability) { + return null; + } + @Override + public void removeNode(NodeInfo node) { + } + @Override + public NodeResponse nodeUpdate(NodeInfo nodeInfo, + Map<CharSequence, List<Container>> containers) { + return null; + } + + @Override + public List<Container> allocate(ApplicationID applicationId, + List<ResourceRequest> ask, List<Container> release) throws IOException { + Container container = new Container(); + container.containerToken = new ContainerToken(); + container.hostName = "localhost"; + container.id = new ContainerID(); + container.id.appID = appID; + container.id.id = count; + count++; + return Arrays.asList(container); + } + + @Override + public void handle(ASMEvent<ApplicationTrackerEventType> event) { + switch (event.getType()) { + case ADD: + schedulerAddApplication++; + break; + case REMOVE: + schedulerRemoveApplication++; + LOG.info("REMOVING app : " + schedulerRemoveApplication); + if (schedulerRemoveApplication == maxFailures) { + synchronized (schedulerNotify) { + schedulerNotify.addAndGet(1); + schedulerNotify.notify(); + } + } + break; + default: + break; + } + } + + @Override + public void reinitialize(Configuration conf, + ContainerTokenSecretManager secretManager) { + } + } + + @Before + public void setUp() { + appID = new ApplicationID(); + appID.clusterTimeStamp = System.currentTimeMillis(); + appID.id = 1; + scheduler = new DummyResourceScheduler(); + asmContext.getDispatcher().register(ApplicationTrackerEventType.class, scheduler); + appImpl = new ExtApplicationsManagerImpl(appTokenSecretManager, scheduler, asmContext); + Configuration conf = new Configuration(); + conf.setLong(YarnConfiguration.AM_EXPIRY_INTERVAL, 1000L); + conf.setInt(YarnConfiguration.AM_MAX_RETRIES, maxFailures); + appImpl.init(conf); + appImpl.start(); + } + + @After + public void tearDown() { + } + + private void waitForFailed(ApplicationMasterInfo masterInfo, ApplicationState + finalState) throws Exception { + int count = 0; + while(masterInfo.getState() != finalState && count < 10) { + Thread.sleep(500); + count++; + } + assertTrue(masterInfo.getState() == finalState); + } + + private class TestAppContext implements AppContext { + private ApplicationID appID; + + public TestAppContext(ApplicationID appID) { + this.appID = appID; + } + @Override + public ApplicationSubmissionContext getSubmissionContext() { + return null; + } + + @Override + public Resource getResource() { + return null; + } + + @Override + public ApplicationID getApplicationID() { + return appID; + } + + @Override + public ApplicationStatus getStatus() { + return null; + } + + @Override + public ApplicationMaster getMaster() { + return null; + } + + @Override + public Container getMasterContainer() { + return null; + } + + @Override + public String getUser() { + return null; + } + + @Override + public long getLastSeen() { + return 0; + } + + @Override + public String getName() { + return null; + } + + @Override + public String getQueue() { + return null; + } + + @Override + public int getFailedCount() { + return 0; + } + + } + + @Test + public void testAMRestart() throws Exception { + ApplicationSubmissionContext subContext = new ApplicationSubmissionContext(); + subContext.applicationId = appID; + subContext.applicationName = "dummyApp"; + subContext.command = new ArrayList<CharSequence>(); + subContext.environment = new HashMap<CharSequence, CharSequence>(); + subContext.fsTokens = new ArrayList<CharSequence>(); + subContext.fsTokens_todo = ByteBuffer.wrap(new byte[0]); + appImpl.submitApplication(subContext); + masterInfo = appImpl.getApplicationMasterInfo(appID); + synchronized (schedulerNotify) { + while(schedulerNotify.get() == 0) { + schedulerNotify.wait(); + } + } + assertTrue(launcherCleanupCalled == maxFailures); + assertTrue(launcherLaunchCalled == maxFailures); + assertTrue(schedulerAddApplication == maxFailures); + assertTrue(schedulerRemoveApplication == maxFailures); + assertTrue(masterInfo.getFailedCount() == maxFailures); + waitForFailed(masterInfo, ApplicationState.FAILED); + stop = true; + } +} \ No newline at end of file Modified: hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestASMStateMachine.java URL: http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestASMStateMachine.java?rev=1084866&r1=1084865&r2=1084866&view=diff ============================================================================== --- hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestASMStateMachine.java (original) +++ hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestASMStateMachine.java Thu Mar 24 07:52:35 2011 @@ -211,6 +211,10 @@ private static class StatusContext imple public String getQueue() { return null; } + @Override + public int getFailedCount() { + return 0; + } } private class ApplicationTracker implements EventHandler<ASMEvent<ApplicationTrackerEventType>> { Modified: hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestApplicationCleanup.java URL: http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestApplicationCleanup.java?rev=1084866&r1=1084865&r2=1084866&view=diff ============================================================================== --- hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestApplicationCleanup.java (original) +++ hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestApplicationCleanup.java Thu Mar 24 07:52:35 2011 @@ -51,9 +51,7 @@ import org.apache.hadoop.yarn.server.res import org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ApplicationMasterEvents.ApplicationTrackerEventType; import org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ApplicationMasterEvents.SNEventType; import org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.NodeInfo; -import org.apache.hadoop.yarn.server.resourcemanager.scheduler.ClusterTracker; -import org.apache.hadoop.yarn.server.resourcemanager.scheduler.ClusterTracker.NodeResponse; -import org.apache.hadoop.yarn.server.resourcemanager.scheduler.ClusterTrackerImpl; +import org.apache.hadoop.yarn.server.resourcemanager.scheduler.NodeResponse; import org.apache.hadoop.yarn.server.resourcemanager.scheduler.ResourceScheduler; import org.apache.hadoop.yarn.server.resourcemanager.scheduler.YarnScheduler; import org.apache.hadoop.yarn.server.resourcemanager.scheduler.fifo.FifoScheduler; @@ -69,25 +67,17 @@ public class TestApplicationCleanup exte private static final Log LOG = LogFactory.getLog(TestApplicationCleanup.class); private AtomicInteger waitForState = new AtomicInteger(0); private ResourceScheduler scheduler; - private ClusterTracker clusterTracker; private final int memoryCapability = 1024; private ExtASM asm; private static final int memoryNeeded = 100; private final ASMContext context = new ResourceManager.ASMContextImpl(); - private class ExtFifoScheduler extends FifoScheduler { - @Override - protected ClusterTracker createClusterTracker() { - clusterTracker = new ClusterTrackerImpl(); - return clusterTracker; - } - } - @Before public void setUp() { new DummyApplicationTracker(); - scheduler = new ExtFifoScheduler(); + scheduler = new FifoScheduler(); + context.getDispatcher().register(ApplicationTrackerEventType.class, scheduler); asm = new ExtASM(new ApplicationTokenSecretManager(), scheduler); asm.init(new Configuration()); } Modified: hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestSchedulerNegotiator.java URL: http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestSchedulerNegotiator.java?rev=1084866&r1=1084865&r2=1084866&view=diff ============================================================================== --- hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestSchedulerNegotiator.java (original) +++ hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestSchedulerNegotiator.java Thu Mar 24 07:52:35 2011 @@ -39,8 +39,9 @@ import org.apache.hadoop.yarn.server.res import org.apache.hadoop.yarn.server.resourcemanager.ResourceManager.ASMContext; import org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ASMEvent; import org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ApplicationMasterEvents.ApplicationEventType; +import org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.events.ApplicationMasterEvents.ApplicationTrackerEventType; import org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.NodeInfo; -import org.apache.hadoop.yarn.server.resourcemanager.scheduler.ClusterTracker.NodeResponse; +import org.apache.hadoop.yarn.server.resourcemanager.scheduler.NodeResponse; import org.apache.hadoop.yarn.server.resourcemanager.scheduler.ResourceScheduler; import org.apache.hadoop.yarn.server.security.ContainerTokenSecretManager; import org.junit.After; @@ -71,15 +72,7 @@ public class TestSchedulerNegotiator ext containers.add(container); return containers; } - @Override - public void addApplication(ApplicationID applicationId, String user, - String unused, Priority priority) - throws IOException { - } - @Override - public void removeApplication(ApplicationID applicationId) - throws IOException { - } + @Override public void reinitialize(Configuration conf, ContainerTokenSecretManager secretManager) { @@ -97,6 +90,10 @@ public class TestSchedulerNegotiator ext @Override public void removeNode(NodeInfo node) { } + + @Override + public void handle(ASMEvent<ApplicationTrackerEventType> event) { + } } @Before Modified: hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/resourcetracker/TestNMExpiry.java URL: http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/resourcetracker/TestNMExpiry.java?rev=1084866&r1=1084865&r2=1084866&view=diff ============================================================================== --- hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/resourcetracker/TestNMExpiry.java (original) +++ hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/resourcetracker/TestNMExpiry.java Thu Mar 24 07:52:35 2011 @@ -33,8 +33,8 @@ import org.apache.hadoop.yarn.conf.YarnC import org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.NodeInfo; import org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.RMResourceTrackerImpl; import org.apache.hadoop.yarn.server.resourcemanager.scheduler.NodeManager; +import org.apache.hadoop.yarn.server.resourcemanager.scheduler.NodeResponse; import org.apache.hadoop.yarn.server.resourcemanager.scheduler.ResourceListener; -import org.apache.hadoop.yarn.server.resourcemanager.scheduler.ClusterTracker.NodeResponse; import org.apache.hadoop.yarn.server.security.ContainerTokenSecretManager; import org.apache.hadoop.yarn.ApplicationID; import org.apache.hadoop.yarn.Container; @@ -74,7 +74,7 @@ public class TestNMExpiry extends TestCa private class TestRMResourceTrackerImpl extends RMResourceTrackerImpl { public TestRMResourceTrackerImpl( ContainerTokenSecretManager containerTokenSecretManager) { - super(containerTokenSecretManager); + super(containerTokenSecretManager, new VoidResourceListener()); } @Override @@ -99,7 +99,6 @@ public class TestNMExpiry extends TestCa @Before public void setUp() { resourceTracker = new TestRMResourceTrackerImpl(containerTokenSecretManager); - resourceTracker.register(new VoidResourceListener()); Configuration conf = new Configuration(); conf.setLong(YarnConfiguration.NM_EXPIRY_INTERVAL, 1000); resourceTracker.init(conf); Modified: hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/resourcetracker/TestRMNMRPCResponseId.java URL: http://svn.apache.org/viewvc/hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/resourcetracker/TestRMNMRPCResponseId.java?rev=1084866&r1=1084865&r2=1084866&view=diff ============================================================================== --- hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/resourcetracker/TestRMNMRPCResponseId.java (original) +++ hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/resourcetracker/TestRMNMRPCResponseId.java Thu Mar 24 07:52:35 2011 @@ -27,8 +27,8 @@ import org.apache.hadoop.net.Node; import org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.NodeInfo; import org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.RMResourceTrackerImpl; import org.apache.hadoop.yarn.server.resourcemanager.scheduler.NodeManager; +import org.apache.hadoop.yarn.server.resourcemanager.scheduler.NodeResponse; import org.apache.hadoop.yarn.server.resourcemanager.scheduler.ResourceListener; -import org.apache.hadoop.yarn.server.resourcemanager.scheduler.ClusterTracker.NodeResponse; import org.apache.hadoop.yarn.server.security.ContainerTokenSecretManager; import org.apache.hadoop.yarn.ApplicationID; import org.apache.hadoop.yarn.Container; @@ -71,8 +71,7 @@ public class TestRMNMRPCResponseId exten @Before public void setUp() { - rmResourceTrackerImpl = new RMResourceTrackerImpl(containerTokenSecretManager); - rmResourceTrackerImpl.register(listener); + rmResourceTrackerImpl = new RMResourceTrackerImpl(containerTokenSecretManager, listener); } @After
