Skip to content

Commit a4c3d48

Browse files
authored
YARN-11504. [Federation] YARN Federation Supports Non-HA mode. (apache#5722)
1 parent 2794fe2 commit a4c3d48

6 files changed

Lines changed: 145 additions & 6 deletions

File tree

hadoop-yarn-project/hadoop-yarn/hadoop-yarn-api/src/main/java/org/apache/hadoop/yarn/conf/YarnConfiguration.java

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3989,6 +3989,10 @@ public static boolean isAclEnabled(Configuration conf) {
39893989
FEDERATION_PREFIX + "failover.enabled";
39903990
public static final boolean DEFAULT_FEDERATION_FAILOVER_ENABLED = true;
39913991

3992+
public static final String FEDERATION_NON_HA_ENABLED =
3993+
FEDERATION_PREFIX + "non-ha.enabled";
3994+
public static final boolean DEFAULT_FEDERATION_NON_HA_ENABLED = false;
3995+
39923996
public static final String FEDERATION_STATESTORE_CLIENT_CLASS =
39933997
FEDERATION_PREFIX + "state-store.class";
39943998

hadoop-yarn-project/hadoop-yarn/hadoop-yarn-common/src/main/java/org/apache/hadoop/yarn/client/AMRMClientUtils.java

Lines changed: 2 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -94,12 +94,8 @@ public static <T> T createRMProxy(final Configuration configuration,
9494
token.setService(ClientRMProxy.getAMRMTokenService(configuration));
9595
setAuthModeInConf(configuration);
9696
}
97-
final T proxyConnection = user.doAs(new PrivilegedExceptionAction<T>() {
98-
@Override
99-
public T run() throws Exception {
100-
return ClientRMProxy.createRMProxy(configuration, protocol);
101-
}
102-
});
97+
final T proxyConnection = user.doAs((PrivilegedExceptionAction<T>) () ->
98+
ClientRMProxy.createRMProxyFederation(configuration, protocol));
10399
return proxyConnection;
104100

105101
} catch (InterruptedException e) {

hadoop-yarn-project/hadoop-yarn/hadoop-yarn-common/src/main/java/org/apache/hadoop/yarn/client/ClientRMProxy.java

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@
2222
import java.net.InetSocketAddress;
2323
import java.util.ArrayList;
2424

25+
import org.apache.hadoop.classification.VisibleForTesting;
2526
import org.slf4j.Logger;
2627
import org.slf4j.LoggerFactory;
2728
import org.apache.hadoop.classification.InterfaceAudience;
@@ -73,6 +74,29 @@ public static <T> T createRMProxy(final Configuration configuration,
7374
return createRMProxy(configuration, protocol, clientRMProxy);
7475
}
7576

77+
/**
78+
* Create a proxy to the ResourceManager for the specified protocol.
79+
* This method is only used for NodeManager#AMRMClientUtils.
80+
*
81+
* @param configuration Configuration with all the required information.
82+
* @param protocol Client protocol for which proxy is being requested.
83+
* @param <T> Type of proxy.
84+
* @return Proxy to the ResourceManager for the specified client protocol.
85+
* @throws IOException io error occur.
86+
*/
87+
public static <T> T createRMProxyFederation(final Configuration configuration,
88+
final Class<T> protocol) throws IOException {
89+
ClientRMProxy<T> clientRMProxy = new ClientRMProxy<>();
90+
return createRMProxyFederation(configuration, protocol, clientRMProxy);
91+
}
92+
93+
@VisibleForTesting
94+
public static <T> RMFailoverProxyProvider<T> getClientRMFailoverProxyProvider(
95+
final YarnConfiguration configuration, final Class<T> protocol) {
96+
ClientRMProxy<T> clientRMProxy = new ClientRMProxy<>();
97+
return getRMFailoverProxyProvider(configuration, protocol, clientRMProxy);
98+
}
99+
76100
private static void setAMRMTokenService(final Configuration conf)
77101
throws IOException {
78102
for (Token<? extends TokenIdentifier> token : UserGroupInformation

hadoop-yarn-project/hadoop-yarn/hadoop-yarn-common/src/main/java/org/apache/hadoop/yarn/client/RMProxy.java

Lines changed: 55 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -116,6 +116,41 @@ protected static <T> T createRMProxy(final Configuration configuration,
116116
return newProxyInstance(conf, protocol, instance, retryPolicy);
117117
}
118118

119+
/**
120+
* This functionality is only used for NodeManager and only in non-HA mode.
121+
* Its purpose is to ensure that when initializes UAM, it can find the correct cluster.
122+
*
123+
* @param configuration configuration.
124+
* @param protocol protocol.
125+
* @param instance RMProxy instance.
126+
* @return RMProxy.
127+
* @param <T> Generic T.
128+
* @throws IOException io error occur.
129+
*/
130+
protected static <T> T createRMProxyFederation(final Configuration configuration,
131+
final Class<T> protocol, RMProxy<T> instance) throws IOException {
132+
YarnConfiguration yarnConf = new YarnConfiguration(configuration);
133+
RetryPolicy retryPolicy = createRetryPolicy(yarnConf, isFailoverEnabled(yarnConf));
134+
return newProxyInstanceFederation(yarnConf, protocol, instance, retryPolicy);
135+
}
136+
137+
protected static <T> T newProxyInstanceFederation(final YarnConfiguration conf,
138+
final Class<T> protocol, RMProxy<T> instance, RetryPolicy retryPolicy) {
139+
RMFailoverProxyProvider<T> provider = getRMFailoverProxyProvider(conf, protocol, instance);
140+
return (T) RetryProxy.create(protocol, provider, retryPolicy);
141+
}
142+
143+
protected static <T> RMFailoverProxyProvider<T> getRMFailoverProxyProvider(
144+
final YarnConfiguration conf, final Class<T> protocol, RMProxy<T> instance) {
145+
RMFailoverProxyProvider<T> provider;
146+
if (isFederationNonHAEnabled(conf)) {
147+
provider = instance.createRMFailoverProxyProvider(conf, protocol);
148+
} else {
149+
provider = instance.createNonHaRMFailoverProxyProvider(conf, protocol);
150+
}
151+
return provider;
152+
}
153+
119154
/**
120155
* Currently, used by NodeManagers only.
121156
* Create a proxy for the specified protocol. For non-HA,
@@ -355,4 +390,24 @@ private static boolean isFailoverEnabled(YarnConfiguration conf) {
355390
return false;
356391
}
357392

393+
/**
394+
* If RM is not configured with HA, NM will not configure yarn.resourcemanager.ha.rmIds locally.
395+
*
396+
* If federation mode is enabled and RMProxy#isFailoverEnabled returns true,
397+
* when NM starts Container, it will try to find the yarn.resourcemanager.ha.rmIds property.
398+
*
399+
* However, an error will occur because this property is not configured
400+
* if the user has not configured HA.
401+
*
402+
* To solve this issue, we can configure the yarn.federation.no-ha.enabled property in NM,
403+
* which tells NM to run in a non-HA environment.
404+
*
405+
* @param conf YarnConfiguration
406+
* @return true, federation support non-HA, false, federation not support non-HA.
407+
*/
408+
private static boolean isFederationNonHAEnabled(YarnConfiguration conf) {
409+
boolean isNonHAEnabled = conf.getBoolean(YarnConfiguration.FEDERATION_NON_HA_ENABLED,
410+
YarnConfiguration.DEFAULT_FEDERATION_NON_HA_ENABLED);
411+
return isNonHAEnabled;
412+
}
358413
}

hadoop-yarn-project/hadoop-yarn/hadoop-yarn-common/src/main/resources/yarn-default.xml

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5354,4 +5354,16 @@
53545354
<value></value>
53555355
</property>
53565356

5357+
<property>
5358+
<description>
5359+
YARN Federation supports Non-HA mode.
5360+
If the cluster is not configured with HA but wants to use YARN Federation,
5361+
this option can be used.
5362+
Setting it to true enables Non-HA mode, while false disables Non-HA mode.
5363+
The default value is false.
5364+
</description>
5365+
<name>yarn.federation.non-ha.enabled</name>
5366+
<value>false</value>
5367+
</property>
5368+
53575369
</configuration>
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,48 @@
1+
/**
2+
* Licensed to the Apache Software Foundation (ASF) under one or more
3+
* contributor license agreements. See the NOTICE file distributed with this
4+
* work for additional information regarding copyright ownership. The ASF
5+
* licenses this file to you under the Apache License, Version 2.0 (the
6+
* "License"); you may not use this file except in compliance with the License.
7+
* You may obtain a copy of the License at
8+
* <p>
9+
* http://www.apache.org/licenses/LICENSE-2.0
10+
* <p>
11+
* Unless required by applicable law or agreed to in writing, software
12+
* distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
13+
* WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
14+
* License for the specific language governing permissions and limitations under
15+
* the License.
16+
*/
17+
package org.apache.hadoop.yarn.server.federation.failover;
18+
19+
import org.apache.hadoop.yarn.api.ApplicationClientProtocol;
20+
import org.apache.hadoop.yarn.client.ClientRMProxy;
21+
import org.apache.hadoop.yarn.client.DefaultNoHARMFailoverProxyProvider;
22+
import org.apache.hadoop.yarn.client.RMFailoverProxyProvider;
23+
import org.apache.hadoop.yarn.conf.YarnConfiguration;
24+
import org.apache.hadoop.yarn.exceptions.YarnException;
25+
import org.junit.Test;
26+
27+
import static org.junit.Assert.assertTrue;
28+
29+
/**
30+
* We will test the failover of Federation.
31+
*/
32+
public class TestFederationRMFailoverProxyProvider {
33+
34+
@Test
35+
public void testRMFailoverProxyProvider() throws YarnException {
36+
YarnConfiguration configuration = new YarnConfiguration();
37+
38+
RMFailoverProxyProvider<ApplicationClientProtocol> clientRMFailoverProxyProvider =
39+
ClientRMProxy.getClientRMFailoverProxyProvider(configuration, ApplicationClientProtocol.class);
40+
assertTrue(clientRMFailoverProxyProvider instanceof DefaultNoHARMFailoverProxyProvider);
41+
42+
FederationProxyProviderUtil.updateConfForFederation(configuration, "SC-1");
43+
configuration.setBoolean(YarnConfiguration.FEDERATION_NON_HA_ENABLED,true);
44+
RMFailoverProxyProvider<ApplicationClientProtocol> clientRMFailoverProxyProvider2 =
45+
ClientRMProxy.getClientRMFailoverProxyProvider(configuration, ApplicationClientProtocol.class);
46+
assertTrue(clientRMFailoverProxyProvider2 instanceof FederationRMFailoverProxyProvider);
47+
}
48+
}

0 commit comments

Comments
 (0)