Repository: hadoop Updated Branches: refs/heads/branch-2 eadd3cecf -> 18c398285
YARN-4227. Ignore expired containers from removed nodes in FairScheduler. (Wilfred Spiegelenburg via rchiang) (cherry picked from commit bc2d67d6c10619716ef7acce263f3269a86c3150) Project: http://git-wip-us.apache.org/repos/asf/hadoop/repo Commit: http://git-wip-us.apache.org/repos/asf/hadoop/commit/18c39828 Tree: http://git-wip-us.apache.org/repos/asf/hadoop/tree/18c39828 Diff: http://git-wip-us.apache.org/repos/asf/hadoop/diff/18c39828 Branch: refs/heads/branch-2 Commit: 18c39828514247994b3eba6a3a9afa0ad23d3b82 Parents: eadd3ce Author: Ray Chiang <[email protected]> Authored: Mon Jan 8 15:32:25 2018 -0800 Committer: Ray Chiang <[email protected]> Committed: Mon Jan 8 16:23:29 2018 -0800 ---------------------------------------------------------------------- .../scheduler/fair/FairScheduler.java | 29 ++++++---- .../scheduler/fair/TestFairScheduler.java | 59 ++++++++++++++++++++ 2 files changed, 78 insertions(+), 10 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/hadoop/blob/18c39828/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/fair/FairScheduler.java ---------------------------------------------------------------------- diff --git a/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/fair/FairScheduler.java b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/fair/FairScheduler.java index a14b2ec..b7fcbe8 100644 --- a/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/fair/FairScheduler.java +++ b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/fair/FairScheduler.java @@ -693,27 +693,36 @@ public class FairScheduler extends ApplicationId appId = container.getId().getApplicationAttemptId().getApplicationId(); if (application == null) { - LOG.info( - "Container " + container + " of" + " finished application " + appId - + " completed with event " + event); + LOG.info("Container " + container + " of finished application " + + appId + " completed with event " + event); return; } // Get the node on which the container was allocated - FSSchedulerNode node = getFSSchedulerNode(container.getNodeId()); - + NodeId nodeID = container.getNodeId(); + FSSchedulerNode node = getFSSchedulerNode(nodeID); + // node could be null if the thread was waiting for the lock and the node + // was removed in another thread if (rmContainer.getState() == RMContainerState.RESERVED) { - application.unreserve(rmContainer.getReservedSchedulerKey(), node); - } else{ + if (node != null) { + application.unreserve(rmContainer.getReservedSchedulerKey(), node); + } else if (LOG.isDebugEnabled()) { + LOG.debug("Skipping unreserve on removed node: " + nodeID); + } + } else { application.containerCompleted(rmContainer, containerStatus, event); - node.releaseContainer(rmContainer.getContainerId(), false); + if (node != null) { + node.releaseContainer(rmContainer.getContainerId(), false); + } else if (LOG.isDebugEnabled()) { + LOG.debug("Skipping container release on removed node: " + nodeID); + } updateRootQueueMetrics(); } if (LOG.isDebugEnabled()) { LOG.debug("Application attempt " + application.getApplicationAttemptId() - + " released container " + container.getId() + " on node: " + node - + " with event: " + event); + + " released container " + container.getId() + " on node: " + + (node == null ? nodeID : node) + " with event: " + event); } } finally { writeLock.unlock(); http://git-wip-us.apache.org/repos/asf/hadoop/blob/18c39828/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/fair/TestFairScheduler.java ---------------------------------------------------------------------- diff --git a/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/fair/TestFairScheduler.java b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/fair/TestFairScheduler.java index a89ba2c..3bc1216 100644 --- a/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/fair/TestFairScheduler.java +++ b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/test/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/fair/TestFairScheduler.java @@ -58,6 +58,7 @@ import org.apache.hadoop.yarn.api.records.ApplicationSubmissionContext; import org.apache.hadoop.yarn.api.records.ContainerId; import org.apache.hadoop.yarn.api.records.ContainerLaunchContext; import org.apache.hadoop.yarn.api.records.FinalApplicationStatus; +import org.apache.hadoop.yarn.api.records.NodeId; import org.apache.hadoop.yarn.api.records.NodeState; import org.apache.hadoop.yarn.api.records.QueueInfo; import org.apache.hadoop.yarn.api.records.Resource; @@ -5353,4 +5354,62 @@ public class TestFairScheduler extends FairSchedulerTestBase { assertTrue(parent.dumpState().equals( parentQueueString + ", " + childQueueString)); } + + @Test + public void testCompletedContainerOnRemovedNode() throws IOException { + scheduler.init(conf); + scheduler.start(); + scheduler.reinitialize(conf, resourceManager.getRMContext()); + + // Add a node + RMNode node = MockNodes.newNodeInfo(1, Resources.createResource(2048), 2, + "127.0.0.2"); + scheduler.handle(new NodeAddedSchedulerEvent(node)); + + // Create application attempt + ApplicationAttemptId appAttemptId = createAppAttemptId(1, 1); + createMockRMApp(appAttemptId); + scheduler.addApplication(appAttemptId.getApplicationId(), "root.queue1", + "user1", false); + scheduler.addApplicationAttempt(appAttemptId, false, false); + + // Create container request that goes to a specific node. + // Without the 2nd and 3rd request we do not get live containers + List<ResourceRequest> ask1 = new ArrayList<>(); + ResourceRequest request1 = + createResourceRequest(1024, node.getHostName(), 1, 1, true); + ask1.add(request1); + ResourceRequest request2 = + createResourceRequest(1024, node.getRackName(), 1, 1, false); + ask1.add(request2); + ResourceRequest request3 = + createResourceRequest(1024, ResourceRequest.ANY, 1, 1, false); + ask1.add(request3); + + // Perform allocation + scheduler.allocate(appAttemptId, ask1, new ArrayList<ContainerId>(), null, + null, NULL_UPDATE_REQUESTS); + scheduler.update(); + scheduler.handle(new NodeUpdateSchedulerEvent(node)); + + // Get the allocated containers for the application (list can not be null) + Collection<RMContainer> clist = scheduler.getSchedulerApp(appAttemptId) + .getLiveContainers(); + Assert.assertEquals(1, clist.size()); + + // Make sure that we remove the correct node (should never fail) + RMContainer rmc = clist.iterator().next(); + NodeId containerNodeID = rmc.getAllocatedNode(); + assertEquals(node.getNodeID(), containerNodeID); + + // Remove node + scheduler.handle(new NodeRemovedSchedulerEvent(node)); + + // Call completedContainer() should not fail even if the node has been + // removed + scheduler.completedContainer(rmc, + SchedulerUtils.createAbnormalContainerStatus(rmc.getContainerId(), + SchedulerUtils.COMPLETED_APPLICATION), + RMContainerEventType.EXPIRE); + } } --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
