YARN-4117. End to end unit test with mini YARN cluster for AMRMProxy Service. Contributed by Giovanni Matteo Fumarola
This commit is contained in:
parent
49ff54c860
commit
55ae143923
@ -28,6 +28,7 @@
|
|||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.concurrent.ConcurrentHashMap;
|
import java.util.concurrent.ConcurrentHashMap;
|
||||||
|
|
||||||
|
import org.apache.hadoop.classification.InterfaceAudience.Private;
|
||||||
import org.apache.hadoop.conf.Configuration;
|
import org.apache.hadoop.conf.Configuration;
|
||||||
import org.apache.hadoop.fs.CommonConfigurationKeysPublic;
|
import org.apache.hadoop.fs.CommonConfigurationKeysPublic;
|
||||||
import org.apache.hadoop.io.DataOutputBuffer;
|
import org.apache.hadoop.io.DataOutputBuffer;
|
||||||
@ -512,6 +513,16 @@ private Token<AMRMTokenIdentifier> getFirstAMRMToken(
|
|||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Private
|
||||||
|
public InetSocketAddress getBindAddress() {
|
||||||
|
return this.listenerEndpoint;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Private
|
||||||
|
public AMRMProxyTokenSecretManager getSecretManager() {
|
||||||
|
return this.secretManager;
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Private class for handling application stop events.
|
* Private class for handling application stop events.
|
||||||
*
|
*
|
||||||
@ -546,7 +557,8 @@ public void handle(ApplicationEvent event) {
|
|||||||
* ApplicationAttemptId instances.
|
* ApplicationAttemptId instances.
|
||||||
*
|
*
|
||||||
*/
|
*/
|
||||||
private static class RequestInterceptorChainWrapper {
|
@Private
|
||||||
|
public static class RequestInterceptorChainWrapper {
|
||||||
private RequestInterceptor rootInterceptor;
|
private RequestInterceptor rootInterceptor;
|
||||||
private ApplicationAttemptId applicationAttemptId;
|
private ApplicationAttemptId applicationAttemptId;
|
||||||
|
|
||||||
|
@ -39,6 +39,8 @@
|
|||||||
import org.slf4j.Logger;
|
import org.slf4j.Logger;
|
||||||
import org.slf4j.LoggerFactory;
|
import org.slf4j.LoggerFactory;
|
||||||
|
|
||||||
|
import com.google.common.annotations.VisibleForTesting;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Extends the AbstractRequestInterceptor class and provides an implementation
|
* Extends the AbstractRequestInterceptor class and provides an implementation
|
||||||
* that simply forwards the AM requests to the cluster resource manager.
|
* that simply forwards the AM requests to the cluster resource manager.
|
||||||
@ -135,4 +137,9 @@ private void updateAMRMToken(Token token) throws IOException {
|
|||||||
user.addToken(amrmToken);
|
user.addToken(amrmToken);
|
||||||
amrmToken.setService(ClientRMProxy.getAMRMTokenService(getConf()));
|
amrmToken.setService(ClientRMProxy.getAMRMTokenService(getConf()));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@VisibleForTesting
|
||||||
|
public void setRMClient(ApplicationMasterProtocol rmClient) {
|
||||||
|
this.rmClient = rmClient;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
@ -183,7 +183,7 @@ public class ContainerManagerImpl extends CompositeService implements
|
|||||||
private final ReadLock readLock;
|
private final ReadLock readLock;
|
||||||
private final WriteLock writeLock;
|
private final WriteLock writeLock;
|
||||||
private AMRMProxyService amrmProxyService;
|
private AMRMProxyService amrmProxyService;
|
||||||
private boolean amrmProxyEnabled = false;
|
protected boolean amrmProxyEnabled = false;
|
||||||
|
|
||||||
private long waitForContainersOnShutdownMillis;
|
private long waitForContainersOnShutdownMillis;
|
||||||
|
|
||||||
@ -247,19 +247,7 @@ public void serviceInit(Configuration conf) throws Exception {
|
|||||||
addService(sharedCacheUploader);
|
addService(sharedCacheUploader);
|
||||||
dispatcher.register(SharedCacheUploadEventType.class, sharedCacheUploader);
|
dispatcher.register(SharedCacheUploadEventType.class, sharedCacheUploader);
|
||||||
|
|
||||||
amrmProxyEnabled =
|
createAMRMProxyService(conf);
|
||||||
conf.getBoolean(YarnConfiguration.AMRM_PROXY_ENABLED,
|
|
||||||
YarnConfiguration.DEFAULT_AMRM_PROXY_ENABLED);
|
|
||||||
|
|
||||||
if (amrmProxyEnabled) {
|
|
||||||
LOG.info("AMRMProxyService is enabled. "
|
|
||||||
+ "All the AM->RM requests will be intercepted by the proxy");
|
|
||||||
this.amrmProxyService =
|
|
||||||
new AMRMProxyService(this.context, this.dispatcher);
|
|
||||||
addService(this.amrmProxyService);
|
|
||||||
} else {
|
|
||||||
LOG.info("AMRMProxyService is disabled");
|
|
||||||
}
|
|
||||||
|
|
||||||
waitForContainersOnShutdownMillis =
|
waitForContainersOnShutdownMillis =
|
||||||
conf.getLong(YarnConfiguration.NM_SLEEP_DELAY_BEFORE_SIGKILL_MS,
|
conf.getLong(YarnConfiguration.NM_SLEEP_DELAY_BEFORE_SIGKILL_MS,
|
||||||
@ -272,8 +260,20 @@ public void serviceInit(Configuration conf) throws Exception {
|
|||||||
recover();
|
recover();
|
||||||
}
|
}
|
||||||
|
|
||||||
public boolean isARMRMProxyEnabled() {
|
protected void createAMRMProxyService(Configuration conf) {
|
||||||
return amrmProxyEnabled;
|
this.amrmProxyEnabled =
|
||||||
|
conf.getBoolean(YarnConfiguration.AMRM_PROXY_ENABLED,
|
||||||
|
YarnConfiguration.DEFAULT_AMRM_PROXY_ENABLED);
|
||||||
|
|
||||||
|
if (amrmProxyEnabled) {
|
||||||
|
LOG.info("AMRMProxyService is enabled. "
|
||||||
|
+ "All the AM->RM requests will be intercepted by the proxy");
|
||||||
|
this.setAMRMProxyService(
|
||||||
|
new AMRMProxyService(this.context, this.dispatcher));
|
||||||
|
addService(this.getAMRMProxyService());
|
||||||
|
} else {
|
||||||
|
LOG.info("AMRMProxyService is disabled");
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@SuppressWarnings("unchecked")
|
@SuppressWarnings("unchecked")
|
||||||
@ -810,9 +810,9 @@ public StartContainersResponse startContainers(
|
|||||||
|
|
||||||
// Initialize the AMRMProxy service instance only if the container is of
|
// Initialize the AMRMProxy service instance only if the container is of
|
||||||
// type AM and if the AMRMProxy service is enabled
|
// type AM and if the AMRMProxy service is enabled
|
||||||
if (isARMRMProxyEnabled() && containerTokenIdentifier
|
if (amrmProxyEnabled && containerTokenIdentifier.getContainerType()
|
||||||
.getContainerType().equals(ContainerType.APPLICATION_MASTER)) {
|
.equals(ContainerType.APPLICATION_MASTER)) {
|
||||||
this.amrmProxyService.processApplicationStartRequest(request);
|
this.getAMRMProxyService().processApplicationStartRequest(request);
|
||||||
}
|
}
|
||||||
|
|
||||||
startContainerInternal(nmTokenIdentifier, containerTokenIdentifier,
|
startContainerInternal(nmTokenIdentifier, containerTokenIdentifier,
|
||||||
@ -1413,4 +1413,15 @@ public Context getContext() {
|
|||||||
public Map<String, ByteBuffer> getAuxServiceMetaData() {
|
public Map<String, ByteBuffer> getAuxServiceMetaData() {
|
||||||
return this.auxiliaryServices.getMetaData();
|
return this.auxiliaryServices.getMetaData();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Private
|
||||||
|
public AMRMProxyService getAMRMProxyService() {
|
||||||
|
return this.amrmProxyService;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Private
|
||||||
|
protected void setAMRMProxyService(AMRMProxyService amrmProxyService) {
|
||||||
|
this.amrmProxyService = amrmProxyService;
|
||||||
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
@ -35,21 +35,23 @@
|
|||||||
import org.apache.hadoop.fs.Path;
|
import org.apache.hadoop.fs.Path;
|
||||||
import org.apache.hadoop.ha.HAServiceProtocol;
|
import org.apache.hadoop.ha.HAServiceProtocol;
|
||||||
import org.apache.hadoop.metrics2.lib.DefaultMetricsSystem;
|
import org.apache.hadoop.metrics2.lib.DefaultMetricsSystem;
|
||||||
|
import org.apache.hadoop.security.token.Token;
|
||||||
import org.apache.hadoop.service.AbstractService;
|
import org.apache.hadoop.service.AbstractService;
|
||||||
import org.apache.hadoop.service.CompositeService;
|
import org.apache.hadoop.service.CompositeService;
|
||||||
import org.apache.hadoop.util.Shell;
|
import org.apache.hadoop.util.Shell;
|
||||||
import org.apache.hadoop.util.Shell.ShellCommandExecutor;
|
import org.apache.hadoop.util.Shell.ShellCommandExecutor;
|
||||||
import org.apache.hadoop.yarn.api.protocolrecords.GetClusterMetricsRequest;
|
import org.apache.hadoop.yarn.api.protocolrecords.GetClusterMetricsRequest;
|
||||||
import org.apache.hadoop.yarn.api.records.ApplicationAttemptId;
|
import org.apache.hadoop.yarn.api.records.ApplicationAttemptId;
|
||||||
import org.apache.hadoop.yarn.api.records.Resource;
|
|
||||||
import org.apache.hadoop.yarn.conf.HAUtil;
|
import org.apache.hadoop.yarn.conf.HAUtil;
|
||||||
import org.apache.hadoop.yarn.conf.YarnConfiguration;
|
import org.apache.hadoop.yarn.conf.YarnConfiguration;
|
||||||
|
import org.apache.hadoop.yarn.event.AsyncDispatcher;
|
||||||
import org.apache.hadoop.yarn.event.Dispatcher;
|
import org.apache.hadoop.yarn.event.Dispatcher;
|
||||||
import org.apache.hadoop.yarn.event.EventHandler;
|
import org.apache.hadoop.yarn.event.EventHandler;
|
||||||
import org.apache.hadoop.yarn.exceptions.YarnException;
|
import org.apache.hadoop.yarn.exceptions.YarnException;
|
||||||
import org.apache.hadoop.yarn.exceptions.YarnRuntimeException;
|
import org.apache.hadoop.yarn.exceptions.YarnRuntimeException;
|
||||||
import org.apache.hadoop.yarn.factories.RecordFactory;
|
import org.apache.hadoop.yarn.factories.RecordFactory;
|
||||||
import org.apache.hadoop.yarn.factory.providers.RecordFactoryProvider;
|
import org.apache.hadoop.yarn.factory.providers.RecordFactoryProvider;
|
||||||
|
import org.apache.hadoop.yarn.security.AMRMTokenIdentifier;
|
||||||
import org.apache.hadoop.yarn.server.api.ResourceTracker;
|
import org.apache.hadoop.yarn.server.api.ResourceTracker;
|
||||||
import org.apache.hadoop.yarn.server.api.protocolrecords.NodeHeartbeatRequest;
|
import org.apache.hadoop.yarn.server.api.protocolrecords.NodeHeartbeatRequest;
|
||||||
import org.apache.hadoop.yarn.server.api.protocolrecords.NodeHeartbeatResponse;
|
import org.apache.hadoop.yarn.server.api.protocolrecords.NodeHeartbeatResponse;
|
||||||
@ -61,24 +63,31 @@
|
|||||||
import org.apache.hadoop.yarn.server.applicationhistoryservice.ApplicationHistoryServer;
|
import org.apache.hadoop.yarn.server.applicationhistoryservice.ApplicationHistoryServer;
|
||||||
import org.apache.hadoop.yarn.server.applicationhistoryservice.ApplicationHistoryStore;
|
import org.apache.hadoop.yarn.server.applicationhistoryservice.ApplicationHistoryStore;
|
||||||
import org.apache.hadoop.yarn.server.applicationhistoryservice.MemoryApplicationHistoryStore;
|
import org.apache.hadoop.yarn.server.applicationhistoryservice.MemoryApplicationHistoryStore;
|
||||||
|
import org.apache.hadoop.yarn.server.nodemanager.ContainerExecutor;
|
||||||
import org.apache.hadoop.yarn.server.nodemanager.Context;
|
import org.apache.hadoop.yarn.server.nodemanager.Context;
|
||||||
|
import org.apache.hadoop.yarn.server.nodemanager.DeletionService;
|
||||||
|
import org.apache.hadoop.yarn.server.nodemanager.LocalDirsHandlerService;
|
||||||
import org.apache.hadoop.yarn.server.nodemanager.NodeHealthCheckerService;
|
import org.apache.hadoop.yarn.server.nodemanager.NodeHealthCheckerService;
|
||||||
import org.apache.hadoop.yarn.server.nodemanager.NodeManager;
|
import org.apache.hadoop.yarn.server.nodemanager.NodeManager;
|
||||||
import org.apache.hadoop.yarn.server.nodemanager.NodeStatusUpdater;
|
import org.apache.hadoop.yarn.server.nodemanager.NodeStatusUpdater;
|
||||||
import org.apache.hadoop.yarn.server.nodemanager.NodeStatusUpdaterImpl;
|
import org.apache.hadoop.yarn.server.nodemanager.NodeStatusUpdaterImpl;
|
||||||
|
import org.apache.hadoop.yarn.server.nodemanager.amrmproxy.AMRMProxyService;
|
||||||
|
import org.apache.hadoop.yarn.server.nodemanager.amrmproxy.DefaultRequestInterceptor;
|
||||||
|
import org.apache.hadoop.yarn.server.nodemanager.amrmproxy.RequestInterceptor;
|
||||||
|
import org.apache.hadoop.yarn.server.nodemanager.containermanager.ContainerManagerImpl;
|
||||||
|
import org.apache.hadoop.yarn.server.nodemanager.metrics.NodeManagerMetrics;
|
||||||
import org.apache.hadoop.yarn.server.resourcemanager.ResourceManager;
|
import org.apache.hadoop.yarn.server.resourcemanager.ResourceManager;
|
||||||
import org.apache.hadoop.yarn.server.resourcemanager.ResourceTrackerService;
|
import org.apache.hadoop.yarn.server.resourcemanager.ResourceTrackerService;
|
||||||
import org.apache.hadoop.yarn.server.resourcemanager.reservation.ReservationSystemTestUtil;
|
|
||||||
import org.apache.hadoop.yarn.server.resourcemanager.rmapp.attempt.RMAppAttemptEvent;
|
import org.apache.hadoop.yarn.server.resourcemanager.rmapp.attempt.RMAppAttemptEvent;
|
||||||
import org.apache.hadoop.yarn.server.resourcemanager.rmapp.attempt.RMAppAttemptEventType;
|
import org.apache.hadoop.yarn.server.resourcemanager.rmapp.attempt.RMAppAttemptEventType;
|
||||||
import org.apache.hadoop.yarn.server.resourcemanager.rmapp.attempt.event.RMAppAttemptRegistrationEvent;
|
import org.apache.hadoop.yarn.server.resourcemanager.rmapp.attempt.event.RMAppAttemptRegistrationEvent;
|
||||||
import org.apache.hadoop.yarn.server.resourcemanager.rmapp.attempt.event.RMAppAttemptUnregistrationEvent;
|
import org.apache.hadoop.yarn.server.resourcemanager.rmapp.attempt.event.RMAppAttemptUnregistrationEvent;
|
||||||
|
import org.apache.hadoop.yarn.server.security.ApplicationACLsManager;
|
||||||
import org.apache.hadoop.yarn.server.timeline.MemoryTimelineStore;
|
import org.apache.hadoop.yarn.server.timeline.MemoryTimelineStore;
|
||||||
import org.apache.hadoop.yarn.server.timeline.TimelineStore;
|
import org.apache.hadoop.yarn.server.timeline.TimelineStore;
|
||||||
import org.apache.hadoop.yarn.server.timeline.recovery.MemoryTimelineStateStore;
|
import org.apache.hadoop.yarn.server.timeline.recovery.MemoryTimelineStateStore;
|
||||||
import org.apache.hadoop.yarn.server.timeline.recovery.TimelineStateStore;
|
import org.apache.hadoop.yarn.server.timeline.recovery.TimelineStateStore;
|
||||||
import org.apache.hadoop.yarn.util.timeline.TimelineUtils;
|
import org.apache.hadoop.yarn.util.timeline.TimelineUtils;
|
||||||
import org.apache.hadoop.yarn.util.resource.Resources;
|
|
||||||
import org.apache.hadoop.yarn.webapp.util.WebAppUtils;
|
import org.apache.hadoop.yarn.webapp.util.WebAppUtils;
|
||||||
|
|
||||||
import com.google.common.annotations.VisibleForTesting;
|
import com.google.common.annotations.VisibleForTesting;
|
||||||
@ -698,6 +707,15 @@ public UnRegisterNodeManagerResponse unRegisterNodeManager(
|
|||||||
protected void stopRMProxy() { }
|
protected void stopRMProxy() { }
|
||||||
};
|
};
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
protected ContainerManagerImpl createContainerManager(Context context,
|
||||||
|
ContainerExecutor exec, DeletionService del,
|
||||||
|
NodeStatusUpdater nodeStatusUpdater, ApplicationACLsManager aclsManager,
|
||||||
|
LocalDirsHandlerService dirsHandler) {
|
||||||
|
return new CustomContainerManagerImpl(context, exec, del,
|
||||||
|
nodeStatusUpdater, metrics, dirsHandler);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@ -799,4 +817,55 @@ protected void doSecureLogin() throws IOException {
|
|||||||
public int getNumOfResourceManager() {
|
public int getNumOfResourceManager() {
|
||||||
return this.resourceManagers.length;
|
return this.resourceManagers.length;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private class CustomContainerManagerImpl extends ContainerManagerImpl {
|
||||||
|
|
||||||
|
public CustomContainerManagerImpl(Context context, ContainerExecutor exec,
|
||||||
|
DeletionService del, NodeStatusUpdater nodeStatusUpdater,
|
||||||
|
NodeManagerMetrics metrics, LocalDirsHandlerService dirsHandler) {
|
||||||
|
super(context, exec, del, nodeStatusUpdater, metrics, dirsHandler);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
protected void createAMRMProxyService(Configuration conf) {
|
||||||
|
this.amrmProxyEnabled =
|
||||||
|
conf.getBoolean(YarnConfiguration.AMRM_PROXY_ENABLED,
|
||||||
|
YarnConfiguration.DEFAULT_AMRM_PROXY_ENABLED);
|
||||||
|
|
||||||
|
if (this.amrmProxyEnabled) {
|
||||||
|
LOG.info("CustomAMRMProxyService is enabled. "
|
||||||
|
+ "All the AM->RM requests will be intercepted by the proxy");
|
||||||
|
AMRMProxyService amrmProxyService =
|
||||||
|
useRpc ? new AMRMProxyService(getContext(), dispatcher)
|
||||||
|
: new ShortCircuitedAMRMProxy(getContext(), dispatcher);
|
||||||
|
this.setAMRMProxyService(amrmProxyService);
|
||||||
|
addService(this.getAMRMProxyService());
|
||||||
|
} else {
|
||||||
|
LOG.info("CustomAMRMProxyService is disabled");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
private class ShortCircuitedAMRMProxy extends AMRMProxyService {
|
||||||
|
|
||||||
|
public ShortCircuitedAMRMProxy(Context context,
|
||||||
|
AsyncDispatcher dispatcher) {
|
||||||
|
super(context, dispatcher);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
protected void initializePipeline(ApplicationAttemptId applicationAttemptId,
|
||||||
|
String user, Token<AMRMTokenIdentifier> amrmToken,
|
||||||
|
Token<AMRMTokenIdentifier> localToken) {
|
||||||
|
super.initializePipeline(applicationAttemptId, user, amrmToken,
|
||||||
|
localToken);
|
||||||
|
RequestInterceptor rt = getPipelines()
|
||||||
|
.get(applicationAttemptId.getApplicationId()).getRootInterceptor();
|
||||||
|
if (rt instanceof DefaultRequestInterceptor) {
|
||||||
|
((DefaultRequestInterceptor) rt)
|
||||||
|
.setRMClient(getResourceManager().getApplicationMasterService());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
Loading…
Reference in New Issue
Block a user