Repository: hadoop Updated Branches: refs/heads/branch-2 fac4172e0 -> 10803bf08
YARN-7173. Container update RM-NM communication fix for backward compatibility. (Arun Suresh via wangda) Change-Id: Ia7d61e0d9df1e703bc983a31e6856e84a5a0521c Project: http://git-wip-us.apache.org/repos/asf/hadoop/repo Commit: http://git-wip-us.apache.org/repos/asf/hadoop/commit/10803bf0 Tree: http://git-wip-us.apache.org/repos/asf/hadoop/tree/10803bf0 Diff: http://git-wip-us.apache.org/repos/asf/hadoop/diff/10803bf0 Branch: refs/heads/branch-2 Commit: 10803bf08dca8af05c0e4006d6892df75210988a Parents: fac4172 Author: Wangda Tan <[email protected]> Authored: Mon Sep 11 20:56:17 2017 -0700 Committer: Wangda Tan <[email protected]> Committed: Mon Sep 11 20:56:17 2017 -0700 ---------------------------------------------------------------------- .../protocolrecords/NodeHeartbeatResponse.java | 5 ++ .../impl/pb/NodeHeartbeatResponsePBImpl.java | 65 ++++++++++++++++++++ .../rmcontainer/RMContainerImpl.java | 4 +- .../resourcemanager/rmnode/RMNodeImpl.java | 20 +++++- .../rmnode/RMNodeUpdateContainerEvent.java | 9 +-- .../scheduler/SchedulerApplicationAttempt.java | 3 +- 6 files changed, 97 insertions(+), 9 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/hadoop/blob/10803bf0/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-common/src/main/java/org/apache/hadoop/yarn/server/api/protocolrecords/NodeHeartbeatResponse.java ---------------------------------------------------------------------- diff --git a/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-common/src/main/java/org/apache/hadoop/yarn/server/api/protocolrecords/NodeHeartbeatResponse.java b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-common/src/main/java/org/apache/hadoop/yarn/server/api/protocolrecords/NodeHeartbeatResponse.java index 4000fc9..8c04ac4 100644 --- a/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-common/src/main/java/org/apache/hadoop/yarn/server/api/protocolrecords/NodeHeartbeatResponse.java +++ b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-common/src/main/java/org/apache/hadoop/yarn/server/api/protocolrecords/NodeHeartbeatResponse.java @@ -89,6 +89,11 @@ public interface NodeHeartbeatResponse { public abstract void addAllContainersToUpdate( Collection<Container> containersToUpdate); + public abstract List<Container> getContainersToDecrease(); + + public abstract void addAllContainersToDecrease( + Collection<Container> containersToDecrease); + ContainerQueuingLimit getContainerQueuingLimit(); void setContainerQueuingLimit(ContainerQueuingLimit containerQueuingLimit); } http://git-wip-us.apache.org/repos/asf/hadoop/blob/10803bf0/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-common/src/main/java/org/apache/hadoop/yarn/server/api/protocolrecords/impl/pb/NodeHeartbeatResponsePBImpl.java ---------------------------------------------------------------------- diff --git a/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-common/src/main/java/org/apache/hadoop/yarn/server/api/protocolrecords/impl/pb/NodeHeartbeatResponsePBImpl.java b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-common/src/main/java/org/apache/hadoop/yarn/server/api/protocolrecords/impl/pb/NodeHeartbeatResponsePBImpl.java index c92e0ea..7ad6b72 100644 --- a/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-common/src/main/java/org/apache/hadoop/yarn/server/api/protocolrecords/impl/pb/NodeHeartbeatResponsePBImpl.java +++ b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-common/src/main/java/org/apache/hadoop/yarn/server/api/protocolrecords/impl/pb/NodeHeartbeatResponsePBImpl.java @@ -73,6 +73,8 @@ public class NodeHeartbeatResponsePBImpl extends private MasterKey nmTokenMasterKey = null; private ContainerQueuingLimit containerQueuingLimit = null; private List<Container> containersToUpdate = null; + // NOTE: This should be removed in 3.x + private List<Container> containersToDecrease = null; private List<SignalContainerRequest> containersToSignal = null; public NodeHeartbeatResponsePBImpl() { @@ -119,6 +121,9 @@ public class NodeHeartbeatResponsePBImpl extends if (this.containersToUpdate != null) { addContainersToUpdateToProto(); } + if (this.containersToDecrease != null) { + addContainersToDecreaseToProto(); + } if (this.containersToSignal != null) { addContainersToSignalToProto(); } @@ -542,6 +547,66 @@ public class NodeHeartbeatResponsePBImpl extends builder.addAllContainersToUpdate(iterable); } + private void initContainersToDecrease() { + if (this.containersToDecrease != null) { + return; + } + NodeHeartbeatResponseProtoOrBuilder p = viaProto ? proto : builder; + List<ContainerProto> list = p.getContainersToDecreaseList(); + this.containersToDecrease = new ArrayList<>(); + + for (ContainerProto c : list) { + this.containersToDecrease.add(convertFromProtoFormat(c)); + } + } + + @Override + public List<Container> getContainersToDecrease() { + initContainersToDecrease(); + return this.containersToDecrease; + } + + @Override + public void addAllContainersToDecrease( + final Collection<Container> containersToDecrease) { + if (containersToDecrease == null) { + return; + } + initContainersToDecrease(); + this.containersToDecrease.addAll(containersToDecrease); + } + + private void addContainersToDecreaseToProto() { + maybeInitBuilder(); + builder.clearContainersToDecrease(); + if (this.containersToDecrease == null) { + return; + } + + Iterable<ContainerProto> iterable = new + Iterable<ContainerProto>() { + @Override + public Iterator<ContainerProto> iterator() { + return new Iterator<ContainerProto>() { + private Iterator<Container> iter = containersToDecrease.iterator(); + @Override + public boolean hasNext() { + return iter.hasNext(); + } + @Override + public ContainerProto next() { + return convertToProtoFormat(iter.next()); + } + @Override + public void remove() { + throw new UnsupportedOperationException(); + } + }; + } + }; + builder.addAllContainersToDecrease(iterable); + } + @Override public Map<ApplicationId, ByteBuffer> getSystemCredentialsForApps() { if (this.systemCredentials != null) { http://git-wip-us.apache.org/repos/asf/hadoop/blob/10803bf0/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/rmcontainer/RMContainerImpl.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/rmcontainer/RMContainerImpl.java b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/rmcontainer/RMContainerImpl.java index 5fee437..0c0b5dd 100644 --- a/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/rmcontainer/RMContainerImpl.java +++ b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/rmcontainer/RMContainerImpl.java @@ -35,6 +35,7 @@ import org.apache.hadoop.yarn.api.records.ContainerId; import org.apache.hadoop.yarn.api.records.ContainerReport; import org.apache.hadoop.yarn.api.records.ContainerState; import org.apache.hadoop.yarn.api.records.ContainerStatus; +import org.apache.hadoop.yarn.api.records.ContainerUpdateType; import org.apache.hadoop.yarn.api.records.ExecutionType; import org.apache.hadoop.yarn.api.records.NodeId; import org.apache.hadoop.yarn.api.records.Priority; @@ -641,7 +642,8 @@ public class RMContainerImpl implements RMContainer { new AllocationExpirationInfo(event.getContainerId())); container.eventHandler.handle(new RMNodeUpdateContainerEvent( container.nodeId, - Collections.singletonList(container.getContainer()))); + Collections.singletonMap(container.getContainer(), + ContainerUpdateType.DECREASE_RESOURCE))); } else if (Resources.fitsIn(nmContainerResource, rmContainerResource)) { // If nmContainerResource < rmContainerResource, this is caused by the // following sequence: http://git-wip-us.apache.org/repos/asf/hadoop/blob/10803bf0/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/rmnode/RMNodeImpl.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/rmnode/RMNodeImpl.java b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/rmnode/RMNodeImpl.java index e6ceae5..83ba8a0 100644 --- a/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/rmnode/RMNodeImpl.java +++ b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/rmnode/RMNodeImpl.java @@ -48,6 +48,7 @@ import org.apache.hadoop.yarn.api.records.Container; import org.apache.hadoop.yarn.api.records.ContainerId; import org.apache.hadoop.yarn.api.records.ContainerState; import org.apache.hadoop.yarn.api.records.ContainerStatus; +import org.apache.hadoop.yarn.api.records.ContainerUpdateType; import org.apache.hadoop.yarn.api.records.ExecutionType; import org.apache.hadoop.yarn.api.records.NodeId; import org.apache.hadoop.yarn.api.records.NodeState; @@ -173,7 +174,11 @@ public class RMNodeImpl implements RMNode, EventHandler<RMNodeEvent> { private final Map<ContainerId, Container> toBeUpdatedContainers = new HashMap<>(); - + + // NOTE: Remove this in 3.x + private final Map<ContainerId, Container> toBeDecreasedContainers = + new HashMap<>(); + private final Map<ContainerId, Container> nmReportedIncreasedContainers = new HashMap<>(); @@ -626,6 +631,10 @@ public class RMNodeImpl implements RMNode, EventHandler<RMNodeEvent> { try { response.addAllContainersToUpdate(toBeUpdatedContainers.values()); toBeUpdatedContainers.clear(); + + // NOTE: Remove in 3.x + response.addAllContainersToDecrease(toBeDecreasedContainers.values()); + toBeDecreasedContainers.clear(); } finally { this.writeLock.unlock(); } @@ -1043,8 +1052,13 @@ public class RMNodeImpl implements RMNode, EventHandler<RMNodeEvent> { public void transition(RMNodeImpl rmNode, RMNodeEvent event) { RMNodeUpdateContainerEvent de = (RMNodeUpdateContainerEvent) event; - for (Container c : de.getToBeUpdatedContainers()) { - rmNode.toBeUpdatedContainers.put(c.getId(), c); + for (Map.Entry<Container, ContainerUpdateType> e : + de.getToBeUpdatedContainers().entrySet()) { + // NOTE: Remove this in 3.x + if (ContainerUpdateType.DECREASE_RESOURCE == e.getValue()) { + rmNode.toBeDecreasedContainers.put(e.getKey().getId(), e.getKey()); + } + rmNode.toBeUpdatedContainers.put(e.getKey().getId(), e.getKey()); } } } http://git-wip-us.apache.org/repos/asf/hadoop/blob/10803bf0/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/rmnode/RMNodeUpdateContainerEvent.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/rmnode/RMNodeUpdateContainerEvent.java b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/rmnode/RMNodeUpdateContainerEvent.java index 73af563..b8f8e73 100644 --- a/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/rmnode/RMNodeUpdateContainerEvent.java +++ b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/rmnode/RMNodeUpdateContainerEvent.java @@ -19,8 +19,10 @@ package org.apache.hadoop.yarn.server.resourcemanager.rmnode; import java.util.List; +import java.util.Map; import org.apache.hadoop.yarn.api.records.Container; +import org.apache.hadoop.yarn.api.records.ContainerUpdateType; import org.apache.hadoop.yarn.api.records.NodeId; /** @@ -29,16 +31,15 @@ import org.apache.hadoop.yarn.api.records.NodeId; * */ public class RMNodeUpdateContainerEvent extends RMNodeEvent { - private List<Container> toBeUpdatedContainers; + private Map<Container, ContainerUpdateType> toBeUpdatedContainers; public RMNodeUpdateContainerEvent(NodeId nodeId, - List<Container> toBeUpdatedContainers) { + Map<Container, ContainerUpdateType> toBeUpdatedContainers) { super(nodeId, RMNodeEventType.UPDATE_CONTAINER); - this.toBeUpdatedContainers = toBeUpdatedContainers; } - public List<Container> getToBeUpdatedContainers() { + public Map<Container, ContainerUpdateType> getToBeUpdatedContainers() { return toBeUpdatedContainers; } } http://git-wip-us.apache.org/repos/asf/hadoop/blob/10803bf0/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/SchedulerApplicationAttempt.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/SchedulerApplicationAttempt.java b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/SchedulerApplicationAttempt.java index 5710fe6..05dc834 100644 --- a/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/SchedulerApplicationAttempt.java +++ b/hadoop-yarn-project/hadoop-yarn/hadoop-yarn-server/hadoop-yarn-server-resourcemanager/src/main/java/org/apache/hadoop/yarn/server/resourcemanager/scheduler/SchedulerApplicationAttempt.java @@ -690,7 +690,8 @@ public class SchedulerApplicationAttempt implements SchedulableEntity { if (autoUpdate) { this.rmContext.getDispatcher().getEventHandler().handle( new RMNodeUpdateContainerEvent(rmContainer.getNodeId(), - Collections.singletonList(rmContainer.getContainer()))); + Collections.singletonMap( + rmContainer.getContainer(), updateType))); } else { rmContainer.handle(new RMContainerUpdatesAcquiredEvent( rmContainer.getContainerId(), --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
