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=1098735&r1=1098734&r2=1098735&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 Mon May 2 19:01:06 2011 @@ -1,20 +1,20 @@ /** -* 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.scheduler.fifo; @@ -52,6 +52,7 @@ import org.apache.hadoop.yarn.security.C 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.resource.Resources; +import org.apache.hadoop.yarn.server.resourcemanager.recovery.ApplicationsStore.ApplicationStore; import org.apache.hadoop.yarn.server.resourcemanager.recovery.Store.ApplicationInfo; import org.apache.hadoop.yarn.server.resourcemanager.recovery.Store.RMState; import org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.ClusterTracker; @@ -67,31 +68,31 @@ import org.apache.hadoop.yarn.server.sec @LimitedPrivate("yarn") @Evolving public class FifoScheduler implements ResourceScheduler { - + private static final Log LOG = LogFactory.getLog(FifoScheduler.class); private final RecordFactory recordFactory = RecordFactoryProvider.getRecordFactory(null); - + Configuration conf; private ContainerTokenSecretManager containerTokenSecretManager; - + // TODO: The memory-block size should be site-configurable? public static final int MINIMUM_MEMORY = 1024; private final static Container[] EMPTY_CONTAINER_ARRAY = new Container[] {}; private final static List<Container> EMPTY_CONTAINER_LIST = Arrays.asList(EMPTY_CONTAINER_ARRAY); private ClusterTracker clusterTracker; - + public static final Resource MINIMUM_ALLOCATION = - Resources.createResource(MINIMUM_MEMORY); - + Resources.createResource(MINIMUM_MEMORY); + Map<ApplicationId, Application> applications = new TreeMap<ApplicationId, Application>( new org.apache.hadoop.yarn.util.BuilderUtils.ApplicationIdComparator()); private static final String DEFAULT_QUEUE_NAME = "default"; private final QueueMetrics metrics = - QueueMetrics.forQueue(DEFAULT_QUEUE_NAME, null, false); + QueueMetrics.forQueue(DEFAULT_QUEUE_NAME, null, false); private final Queue DEFAULT_QUEUE = new Queue() { @Override @@ -145,14 +146,14 @@ public class FifoScheduler implements Re public FifoScheduler() { } - + public FifoScheduler(ClusterTracker clusterTracker) { this.clusterTracker = clusterTracker; if (clusterTracker != null) { this.clusterTracker.addListener(this); } } - + @Override public void reinitialize(Configuration conf, ContainerTokenSecretManager containerTokenSecretManager, ClusterTracker clusterTracker) @@ -163,7 +164,7 @@ public class FifoScheduler implements Re this.clusterTracker = clusterTracker; if (clusterTracker != null) this.clusterTracker.addListener(this); } - + @Override public synchronized List<Container> allocate(ApplicationId applicationId, List<ResourceRequest> ask, List<Container> release) @@ -171,19 +172,19 @@ public class FifoScheduler implements Re Application application = getApplication(applicationId); if (application == null) { LOG.error("Calling allocate on removed " + - "or non existant application " + applicationId); + "or non existant application " + applicationId); return EMPTY_CONTAINER_LIST; } normalizeRequests(ask); - + LOG.debug("allocate: pre-update" + - " applicationId=" + applicationId + + " applicationId=" + applicationId + " application=" + application); application.showRequests(); - + // Update application requests application.updateResourceRequests(ask); - + // Release containers releaseContainers(application, release); @@ -191,13 +192,13 @@ public class FifoScheduler implements Re " applicationId=" + applicationId + " application=" + application); application.showRequests(); - + List<Container> allContainers = application.acquire(); LOG.debug("allocate:" + - " applicationId=" + applicationId + - " #ask=" + ask.size() + - " #release=" + release.size() + - " #allContainers=" + allContainers.size()); + " applicationId=" + applicationId + + " #ask=" + ask.size() + + " #release=" + release.size() + + " #allContainers=" + allContainers.size()); return allContainers; } @@ -207,31 +208,31 @@ public class FifoScheduler implements Re releaseContainer(application.getApplicationId(), container); } } - + private void normalizeRequests(List<ResourceRequest> asks) { for (ResourceRequest ask : asks) { normalizeRequest(ask); } } - + private void normalizeRequest(ResourceRequest ask) { int memory = ask.getCapability().getMemory(); // FIXME: TestApplicationCleanup is relying on unnormalized behavior. //ask.capability.memory = MINIMUM_MEMORY * memory = MINIMUM_MEMORY * - ((memory/MINIMUM_MEMORY) + (memory%MINIMUM_MEMORY > 0 ? 1 : 0)); + ((memory/MINIMUM_MEMORY) + (memory%MINIMUM_MEMORY > 0 ? 1 : 0)); } - + private synchronized Application getApplication(ApplicationId applicationId) { return applications.get(applicationId); } @Override public synchronized void addApplication(ApplicationId applicationId, ApplicationMaster master, - String user, String unusedQueue, Priority unusedPriority) + String user, String unusedQueue, Priority unusedPriority, ApplicationStore appStore) throws IOException { applications.put(applicationId, - new Application(applicationId, master, DEFAULT_QUEUE, user)); + new Application(applicationId, master, DEFAULT_QUEUE, user, appStore)); metrics.submitApp(user); LOG.info("Application Submission: " + applicationId.getId() + " from " + user + ", currently active: " + applications.size()); @@ -243,9 +244,9 @@ public class FifoScheduler implements Re Application application = getApplication(applicationId); if (application == null) { throw new IOException("Unknown application " + applicationId + - " has completed!"); + " has completed!"); } - + // Release current containers releaseContainers(application, application.getCurrentContainers()); @@ -260,7 +261,7 @@ public class FifoScheduler implements Re // Remove the application applications.remove(applicationId); } - + /** * Heart of the scheduler... * @@ -268,9 +269,9 @@ public class FifoScheduler implements Re */ private synchronized void assignContainers(NodeInfo node) { LOG.debug("assignContainers:" + - " node=" + node.getNodeAddress() + + " node=" + node.getNodeAddress() + " #applications=" + applications.size()); - + // Try to assign containers to applications in fifo order for (Map.Entry<ApplicationId, Application> e : applications.entrySet()) { Application application = e.getValue(); @@ -294,24 +295,24 @@ public class FifoScheduler implements Re } LOG.debug("post-assignContainers"); application.showRequests(); - + // Done if (Resources.lessThan(node.getAvailableResource(), MINIMUM_ALLOCATION)) { return; } } } - + private int getMaxAllocatableContainers(Application application, Priority priority, NodeInfo node, NodeType type) { ResourceRequest offSwitchRequest = application.getResourceRequest(priority, NodeManager.ANY); int maxContainers = offSwitchRequest.getNumContainers(); - + if (type == NodeType.OFF_SWITCH) { return maxContainers; } - + if (type == NodeType.RACK_LOCAL) { ResourceRequest rackLocalRequest = application.getResourceRequest(priority, node.getRackName()); @@ -321,7 +322,7 @@ public class FifoScheduler implements Re maxContainers = Math.min(maxContainers, rackLocalRequest.getNumContainers()); } - + if (type == NodeType.DATA_LOCAL) { ResourceRequest nodeLocalRequest = application.getResourceRequest(priority, node.getNodeAddress()); @@ -329,14 +330,14 @@ public class FifoScheduler implements Re maxContainers = Math.min(maxContainers, nodeLocalRequest.getNumContainers()); } } - + return maxContainers; } - + private int assignContainersOnNode(NodeInfo node, Application application, Priority priority - ) { + ) { // Data-local int nodeLocalContainers = assignNodeLocalContainers(node, application, priority); @@ -344,23 +345,23 @@ public class FifoScheduler implements Re // Rack-local int rackLocalContainers = assignRackLocalContainers(node, application, priority); - + // Off-switch int offSwitchContainers = assignOffSwitchContainers(node, application, priority); - + LOG.debug("assignContainersOnNode:" + " node=" + node.getNodeAddress() + " application=" + application.getApplicationId().getId() + " priority=" + priority.getPriority() + " #assigned=" + - (nodeLocalContainers + rackLocalContainers + offSwitchContainers)); - + (nodeLocalContainers + rackLocalContainers + offSwitchContainers)); + return (nodeLocalContainers + rackLocalContainers + offSwitchContainers); } - + private int assignNodeLocalContainers(NodeInfo node, Application application, Priority priority) { int assignedContainers = 0; @@ -371,14 +372,14 @@ public class FifoScheduler implements Re Math.min( getMaxAllocatableContainers(application, priority, node, NodeType.DATA_LOCAL), - request.getNumContainers()); + request.getNumContainers()); assignedContainers = assignContainers(node, application, priority, assignableContainers, request, NodeType.DATA_LOCAL); } return assignedContainers; } - + private int assignRackLocalContainers(NodeInfo node, Application application, Priority priority) { int assignedContainers = 0; @@ -389,7 +390,7 @@ public class FifoScheduler implements Re Math.min( getMaxAllocatableContainers(application, priority, node, NodeType.RACK_LOCAL), - request.getNumContainers()); + request.getNumContainers()); assignedContainers = assignContainers(node, application, priority, assignableContainers, request, NodeType.RACK_LOCAL); @@ -409,43 +410,43 @@ public class FifoScheduler implements Re } return assignedContainers; } - + private int assignContainers(NodeInfo node, Application application, Priority priority, int assignableContainers, ResourceRequest request, NodeType type) { LOG.debug("assignContainers:" + - " node=" + node.getNodeAddress() + - " application=" + application.getApplicationId().getId() + + " node=" + node.getNodeAddress() + + " application=" + application.getApplicationId().getId() + " priority=" + priority.getPriority() + " assignableContainers=" + assignableContainers + " request=" + request + " type=" + type); Resource capability = request.getCapability(); - + int availableContainers = - node.getAvailableResource().getMemory() / capability.getMemory(); // TODO: A buggy - // application - // with this - // zero would - // crash the - // scheduler. + node.getAvailableResource().getMemory() / capability.getMemory(); // TODO: A buggy + // application + // with this + // zero would + // crash the + // scheduler. int assignedContainers = Math.min(assignableContainers, availableContainers); - + if (assignedContainers > 0) { List<Container> containers = - new ArrayList<Container>(assignedContainers); + new ArrayList<Container>(assignedContainers); for (int i=0; i < assignedContainers; ++i) { Container container = - org.apache.hadoop.yarn.server.resourcemanager.resource.Container - .create(recordFactory, application.getApplicationId(), - application.getNewContainerId(), node.getNodeAddress(), - node.getHttpAddress(), capability); + org.apache.hadoop.yarn.server.resourcemanager.resource.Container + .create(recordFactory, application.getApplicationId(), + application.getNewContainerId(), node.getNodeAddress(), + node.getHttpAddress(), capability); // If security is enabled, send the container-tokens too. if (UserGroupInformation.isSecurityEnabled()) { ContainerToken containerToken = recordFactory.newRecordInstance(ContainerToken.class); ContainerTokenIdentifier tokenidentifier = - new ContainerTokenIdentifier(container.getId(), - container.getContainerManagerAddress(), container.getResource()); + new ContainerTokenIdentifier(container.getId(), + container.getContainerManagerAddress(), container.getResource()); containerToken.setIdentifier( ByteBuffer.wrap(tokenidentifier.getBytes())); containerToken.setKind(ContainerTokenIdentifier.KIND.toString()); @@ -460,13 +461,13 @@ public class FifoScheduler implements Re application.allocate(type, node, priority, request, containers); addAllocatedContainers(node, application.getApplicationId(), containers); Resources.addTo(usedResource, - Resources.multiply(capability, assignedContainers)); + Resources.multiply(capability, assignedContainers)); } return assignedContainers; } private synchronized void applicationCompletedContainers( - List<Container> completedContainers) { + List<Container> completedContainers) { for (Container c: completedContainers) { Application app = applications.get(c.getId().getAppId()); /** this is possible, since an application can be removed from scheduler but @@ -477,7 +478,7 @@ public class FifoScheduler implements Re } } } - + private List<Container> getCompletedContainers(Map<String, List<Container>> allContainers) { if (allContainers == null) { return new ArrayList<Container>(); @@ -494,21 +495,21 @@ public class FifoScheduler implements Re } return completedContainers; } - + @Override public synchronized void nodeUpdate(NodeInfo node, Map<String,List<Container>> containers ) { - + applicationCompletedContainers(getCompletedContainers(containers)); LOG.info("Node heartbeat " + node.getNodeID() + " resource = " + node.getAvailableResource()); if (Resources.greaterThanOrEqual(node.getAvailableResource(), - MINIMUM_ALLOCATION)) { + MINIMUM_ALLOCATION)) { assignContainers(node); } metrics.setAvailableQueueMemory( clusterResource.getMemory() - usedResource.getMemory()); LOG.info("Node after allocation " + node.getNodeID() + " resource = " - + node.getAvailableResource()); + + node.getAvailableResource()); // TODO: Add the list of containers to be preempted when we support } @@ -520,7 +521,8 @@ public class FifoScheduler implements Re try { addApplication(event.getAppContext().getApplicationID(), event.getAppContext().getMaster(), event.getAppContext().getUser(), - event.getAppContext().getQueue(), event.getAppContext().getSubmissionContext().getPriority()); + event.getAppContext().getQueue(), event.getAppContext().getSubmissionContext().getPriority() + , event.getAppContext().getStore()); } 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 **/ @@ -543,12 +545,12 @@ public class FifoScheduler implements Re break; } } - + private Map<String, NodeManager> nodes = new HashMap<String, NodeManager>(); private Resource clusterResource = recordFactory.newRecordInstance(Resource.class); private Resource usedResource = recordFactory.newRecordInstance(Resource.class); - + public synchronized Resource getClusterResource() { return clusterResource; } @@ -558,14 +560,14 @@ public class FifoScheduler implements Re Resources.subtractFrom(clusterResource, nodeInfo.getTotalCapability()); //TODO inform the the applications that the containers are completed/failed } - + @Override public QueueInfo getQueueInfo(String queueName, boolean includeApplications, boolean includeChildQueues, boolean recursive) { return DEFAULT_QUEUE.getQueueInfo(includeApplications, false, false); } - + @Override public List<QueueUserACLInfo> getQueueUserAclInfo() { return DEFAULT_QUEUE.getQueueUserAclInfo(null); @@ -585,14 +587,14 @@ public class FifoScheduler implements Re public synchronized void addNode(NodeInfo nodeManager) { Resources.addTo(clusterResource, nodeManager.getTotalCapability()); } - + public synchronized void releaseContainer(ApplicationId applicationId, Container container) { // Reap containers LOG.info("Application " + applicationId + " released container " + container); clusterTracker.releaseContainer(container); } - + public synchronized void addAllocatedContainers(NodeInfo nodeInfo, ApplicationId applicationId, List<Container> containers) { nodeInfo.allocateContainer(applicationId, containers); @@ -605,8 +607,11 @@ public class FifoScheduler implements Re @Override public void recover(RMState state) { - /** nothing special to do - * ASM recovery should take care of the scheduler recovery. - */ + for (Map.Entry<ApplicationId, ApplicationInfo> entry: state.getStoredApplications().entrySet()) { + ApplicationId appId = entry.getKey(); + ApplicationInfo appInfo = entry.getValue(); + Application app = applications.get(appId); + app.allocate(appInfo.getContainers()); + } } }
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=1098735&r1=1098734&r2=1098735&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 Mon May 2 19:01:06 2011 @@ -48,6 +48,7 @@ import org.apache.hadoop.yarn.server.res 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.recovery.ApplicationsStore.ApplicationStore; import org.apache.hadoop.yarn.server.resourcemanager.recovery.MemStore; import org.apache.hadoop.yarn.server.resourcemanager.scheduler.YarnScheduler; import org.junit.After; @@ -93,7 +94,8 @@ public class TestAMLaunchFailure extends @Override public void addApplication(ApplicationId applicationId, - ApplicationMaster master, String user, String queue, Priority priority) + ApplicationMaster master, String user, String queue, Priority priority + , ApplicationStore appStore) throws IOException { // TODO Auto-generated method stub 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=1098735&r1=1098734&r2=1098735&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 Mon May 2 19:01:06 2011 @@ -43,6 +43,8 @@ import org.apache.hadoop.yarn.server.res import org.apache.hadoop.yarn.server.resourcemanager.ResourceManager; import org.apache.hadoop.yarn.server.resourcemanager.ResourceManager.RMContext; import org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager.ApplicationsManagerImpl; +import org.apache.hadoop.yarn.server.resourcemanager.recovery.ApplicationsStore; +import org.apache.hadoop.yarn.server.resourcemanager.recovery.ApplicationsStore.ApplicationStore; import org.apache.hadoop.yarn.server.resourcemanager.recovery.MemStore; import org.apache.hadoop.yarn.server.resourcemanager.scheduler.YarnScheduler; import org.junit.After; @@ -103,10 +105,9 @@ public class TestAMRMRPCResponseId exten @Override public void addApplication(ApplicationId applicationId, - ApplicationMaster master, String user, String queue, Priority priority) - throws IOException { - // TODO Auto-generated method stub - + ApplicationMaster master, String user, String queue, Priority priority, + ApplicationStore store) + throws IOException { } } Modified: 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=1098735&r1=1098734&r2=1098735&view=diff ============================================================================== --- hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMRestart.java (original) +++ hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestAMRestart.java Mon May 2 19:01:06 2011 @@ -36,7 +36,9 @@ import org.apache.hadoop.yarn.server.res 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.recovery.ApplicationsStore.ApplicationStore; import org.apache.hadoop.yarn.server.resourcemanager.recovery.MemStore; +import org.apache.hadoop.yarn.server.resourcemanager.recovery.StoreFactory; import org.apache.hadoop.yarn.server.resourcemanager.recovery.Store.RMState; import org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.ClusterTracker; import org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.NodeInfo; @@ -187,7 +189,8 @@ public class TestAMRestart extends TestC } @Override public void addApplication(ApplicationId applicationId, - ApplicationMaster master, String user, String queue, Priority priority) + ApplicationMaster master, String user, String queue, Priority priority, + ApplicationStore store) throws IOException { } @Override @@ -301,6 +304,11 @@ public class TestAMRestart extends TestC return 0; } + @Override + public ApplicationStore getStore() { + return StoreFactory.createVoidAppStore(); + } + } @Test 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=1098735&r1=1098734&r2=1098735&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 Mon May 2 19:01:06 2011 @@ -43,7 +43,9 @@ import org.apache.hadoop.yarn.server.res 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.applicationsmanager.events.ApplicationMasterEvents.SNEventType; +import org.apache.hadoop.yarn.server.resourcemanager.recovery.ApplicationsStore.ApplicationStore; import org.apache.hadoop.yarn.server.resourcemanager.recovery.MemStore; +import org.apache.hadoop.yarn.server.resourcemanager.recovery.StoreFactory; import org.junit.After; import org.junit.Before; import org.junit.Test; @@ -175,6 +177,10 @@ public class TestASMStateMachine extends public int getFailedCount() { return 0; } + @Override + public ApplicationStore getStore() { + return StoreFactory.createVoidAppStore(); + } } private class ApplicationTracker implements EventHandler<ASMEvent<ApplicationTrackerEventType>> { @@ -227,7 +233,8 @@ public class TestASMStateMachine extends submissioncontext.getApplicationId().setClusterTimestamp(System.currentTimeMillis()); ApplicationMasterInfo masterInfo - = new ApplicationMasterInfo(context, "dummyuser", submissioncontext, "dummyToken"); + = new ApplicationMasterInfo(context, "dummyuser", submissioncontext, "dummyToken" + , StoreFactory.createVoidAppStore()); context.getDispatcher().register(ApplicationEventType.class, masterInfo); handler.handle(new ASMEvent<ApplicationEventType>(ApplicationEventType. 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=1098735&r1=1098734&r2=1098735&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 Mon May 2 19:01:06 2011 @@ -18,6 +18,7 @@ package org.apache.hadoop.yarn.server.resourcemanager.applicationsmanager; +import java.io.IOException; import java.util.ArrayList; import java.util.HashMap; import java.util.List; @@ -57,6 +58,7 @@ import org.apache.hadoop.yarn.server.res 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.applicationsmanager.events.ApplicationMasterEvents.SNEventType; +import org.apache.hadoop.yarn.server.resourcemanager.recovery.ApplicationsStore.ApplicationStore; import org.apache.hadoop.yarn.server.resourcemanager.recovery.MemStore; import org.apache.hadoop.yarn.server.resourcemanager.recovery.Store.RMState; import org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.ClusterTracker; @@ -94,13 +96,13 @@ public class TestApplicationCleanup exte private class TestRMResourceTrackerImpl extends RMResourceTrackerImpl { public TestRMResourceTrackerImpl( - ContainerTokenSecretManager containerTokenSecretManager) { - super(containerTokenSecretManager); + ContainerTokenSecretManager containerTokenSecretManager, RMContext context) { + super(containerTokenSecretManager, context); } @Override protected NodeInfoTracker getAndAddNodeInfoTracker(NodeId nodeId, - String hostString, String httpAddress, Node node, Resource capability) { + String hostString, String httpAddress, Node node, Resource capability) throws IOException { return super.getAndAddNodeInfoTracker(nodeId, hostString, httpAddress, node, capability); } @@ -113,7 +115,7 @@ public class TestApplicationCleanup exte @Before public void setUp() { new DummyApplicationTracker(); - this.clusterTracker = new TestRMResourceTrackerImpl(new ContainerTokenSecretManager()); + this.clusterTracker = new TestRMResourceTrackerImpl(new ContainerTokenSecretManager(), context); scheduler = new FifoScheduler(clusterTracker); context.getDispatcher().register(ApplicationTrackerEventType.class, scheduler); Configuration conf = new Configuration(); @@ -242,7 +244,7 @@ public class TestApplicationCleanup exte return request; } - protected NodeManager addNodes(String commonName, int i, int memoryCapability) { + protected NodeManager addNodes(String commonName, int i, int memoryCapability) throws IOException { NodeId nodeId = recordFactory.newRecordInstance(NodeId.class); nodeId.setId(i); String hostName = commonName + "_" + i; Modified: hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestApplicationMasterLauncher.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/TestApplicationMasterLauncher.java?rev=1098735&r1=1098734&r2=1098735&view=diff ============================================================================== --- hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestApplicationMasterLauncher.java (original) +++ hadoop/mapreduce/branches/MR-279/yarn/yarn-server/yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/applicationsmanager/TestApplicationMasterLauncher.java Mon May 2 19:01:06 2011 @@ -38,6 +38,7 @@ import org.apache.hadoop.yarn.server.res 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.recovery.MemStore; +import org.apache.hadoop.yarn.server.resourcemanager.recovery.StoreFactory; import org.junit.After; import org.junit.Before; import org.junit.Test; @@ -150,7 +151,7 @@ public class TestApplicationMasterLaunch context.setUser("dummyuser"); ApplicationMasterInfo masterInfo = new ApplicationMasterInfo(this.context, "dummyuser", context, - "dummyclienttoken"); + "dummyclienttoken", StoreFactory.createVoidAppStore()); amLauncher.handle(new ASMEvent<AMLauncherEventType>(AMLauncherEventType.LAUNCH, masterInfo)); amLauncher.handle(new ASMEvent<AMLauncherEventType>(AMLauncherEventType.CLEANUP, 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=1098735&r1=1098734&r2=1098735&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 Mon May 2 19:01:06 2011 @@ -45,9 +45,11 @@ import org.apache.hadoop.yarn.server.res 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.recovery.ApplicationsStore.ApplicationStore; import org.apache.hadoop.yarn.server.resourcemanager.recovery.MemStore; import org.apache.hadoop.yarn.server.resourcemanager.recovery.Store; import org.apache.hadoop.yarn.server.resourcemanager.recovery.Store.RMState; +import org.apache.hadoop.yarn.server.resourcemanager.recovery.StoreFactory; import org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.ClusterTracker; import org.apache.hadoop.yarn.server.resourcemanager.resourcetracker.NodeInfo; import org.apache.hadoop.yarn.server.resourcemanager.scheduler.NodeManager; @@ -111,7 +113,8 @@ public class TestSchedulerNegotiator ext } @Override public void addApplication(ApplicationId applicationId, - ApplicationMaster master, String user, String queue, Priority priority) + ApplicationMaster master, String user, String queue, Priority priority, + ApplicationStore store) throws IOException { } @@ -178,7 +181,7 @@ public class TestSchedulerNegotiator ext masterInfo = new ApplicationMasterInfo(this.context, - "dummy", submissionContext, "dummyClientToken"); + "dummy", submissionContext, "dummyClientToken", StoreFactory.createVoidAppStore()); context.getDispatcher().register(ApplicationEventType.class, masterInfo); context.getDispatcher().register(ApplicationTrackerEventType.class, masterInfo); handler.handle(new ASMEvent<ApplicationEventType>(ApplicationEventType. 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=1098735&r1=1098734&r2=1098735&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 Mon May 2 19:01:06 2011 @@ -40,6 +40,10 @@ import org.apache.hadoop.yarn.server.api import org.apache.hadoop.yarn.server.api.protocolrecords.RegisterNodeManagerRequest; import org.apache.hadoop.yarn.server.api.records.HeartbeatResponse; import org.apache.hadoop.yarn.server.api.records.RegistrationResponse; +import org.apache.hadoop.yarn.server.resourcemanager.ResourceManager; +import org.apache.hadoop.yarn.server.resourcemanager.ResourceManager.RMContext; +import org.apache.hadoop.yarn.server.resourcemanager.recovery.MemStore; +import org.apache.hadoop.yarn.server.resourcemanager.recovery.StoreFactory; 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; @@ -77,8 +81,8 @@ public class TestNMExpiry extends TestCa private class TestRMResourceTrackerImpl extends RMResourceTrackerImpl { public TestRMResourceTrackerImpl( - ContainerTokenSecretManager containerTokenSecretManager) { - super(containerTokenSecretManager); + ContainerTokenSecretManager containerTokenSecretManager, RMContext context) { + super(containerTokenSecretManager, context); } @Override @@ -103,9 +107,12 @@ public class TestNMExpiry extends TestCa @Before public void setUp() { - resourceTracker = new TestRMResourceTrackerImpl(containerTokenSecretManager); - resourceTracker.addListener(new VoidResourceListener()); Configuration conf = new Configuration(); + RMContext context = new ResourceManager.RMContextImpl(new MemStore()); + resourceTracker = new TestRMResourceTrackerImpl(containerTokenSecretManager, + context); + resourceTracker.addListener(new VoidResourceListener()); + conf.setLong(YarnConfiguration.NM_EXPIRY_INTERVAL, 1000); resourceTracker.init(conf); resourceTracker.start(); 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=1098735&r1=1098734&r2=1098735&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 Mon May 2 19:01:06 2011 @@ -33,6 +33,9 @@ import org.apache.hadoop.yarn.factory.pr import org.apache.hadoop.yarn.server.api.protocolrecords.NodeHeartbeatRequest; import org.apache.hadoop.yarn.server.api.protocolrecords.RegisterNodeManagerRequest; import org.apache.hadoop.yarn.server.api.records.HeartbeatResponse; +import org.apache.hadoop.yarn.server.resourcemanager.ResourceManager; +import org.apache.hadoop.yarn.server.resourcemanager.ResourceManager.RMContext; +import org.apache.hadoop.yarn.server.resourcemanager.recovery.MemStore; 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; @@ -73,7 +76,8 @@ public class TestRMNMRPCResponseId exten @Before public void setUp() { - rmResourceTrackerImpl = new RMResourceTrackerImpl(containerTokenSecretManager); + RMContext context = new ResourceManager.RMContextImpl(new MemStore()); + rmResourceTrackerImpl = new RMResourceTrackerImpl(containerTokenSecretManager, context); rmResourceTrackerImpl.addListener(listener); }
