This is an automated email from the ASF dual-hosted git repository.
zhouky pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/celeborn.git
The following commit(s) were added to refs/heads/main by this push:
new b7fa30622 [CELEBORN-1354][FOLLOWUP] Split rpc_app into
rpc_app_lifecyclemanager and rpc_app_client
b7fa30622 is described below
commit b7fa30622e1fd4e2d6fd1c2c1e49e4aa1b30b1dd
Author: Mridul Muralidharan <[email protected]>
AuthorDate: Mon Apr 22 19:57:10 2024 +0800
[CELEBORN-1354][FOLLOWUP] Split rpc_app into rpc_app_lifecyclemanager and
rpc_app_client
### What changes were proposed in this pull request?
Based on [review feedback
here](https://github.com/apache/celeborn/pull/2469#discussion_r1573228448),
split `rpc_app` module into two internal transport modules -
`rpc_app_lifecyclemanager` and `rpc_app_client`.
These are not directly configured by users - but internally allow us to
differentiate between lifecycle manager and executor
### Why are the changes needed?
auto ssl is to be configured only at lifecycle manager - doing it at
executors is beneign, but not very useful (it results in creation of local
files, etc).
Avoid this by handling it specifically only for lifecyclemanager.
This is based on [review
feedback](https://github.com/apache/celeborn/pull/2469#discussion_r1573228448).
### Does this PR introduce _any_ user-facing change?
No, the modules are all internal and not user visible.
### How was this patch tested?
Unit tests were updated
Closes #2471 from mridulm/support-app-auto-ssl-followup.
Lead-authored-by: Mridul Muralidharan <[email protected]>
Co-authored-by: Mridul Muralidharan <mridulatgmail.com>
Signed-off-by: zky.zhoukeyong <[email protected]>
---
.../apache/celeborn/client/ShuffleClientImpl.java | 2 +-
.../apache/celeborn/client/LifecycleManager.scala | 2 +-
.../celeborn/common/network/TransportContext.java | 9 +++++
.../common/network/util/TransportConf.java | 4 +-
.../common/protocol/TransportModuleConstants.java | 10 +++++
.../org/apache/celeborn/common/CelebornConf.scala | 43 ++++++++++++++++++++--
.../org/apache/celeborn/common/rpc/RpcEnv.scala | 13 +++++--
.../org/apache/celeborn/common/util/Utils.scala | 15 ++++++--
.../common/network/AutoSSLRpcIntegrationSuite.java | 10 +++--
.../common/network/ssl/SslConnectivitySuiteJ.java | 38 ++++++++++++++-----
.../common/network/util/TransportConfSuiteJ.java | 2 +-
.../apache/celeborn/common/rpc/RpcEnvSuite.scala | 4 +-
12 files changed, 124 insertions(+), 28 deletions(-)
diff --git
a/client/src/main/java/org/apache/celeborn/client/ShuffleClientImpl.java
b/client/src/main/java/org/apache/celeborn/client/ShuffleClientImpl.java
index 3be2e579e..11f13ea36 100644
--- a/client/src/main/java/org/apache/celeborn/client/ShuffleClientImpl.java
+++ b/client/src/main/java/org/apache/celeborn/client/ShuffleClientImpl.java
@@ -190,7 +190,7 @@ public class ShuffleClientImpl extends ShuffleClient {
rpcEnv =
RpcEnv.create(
RpcNameConstants.SHUFFLE_CLIENT_SYS,
- TransportModuleConstants.RPC_APP_MODULE,
+ TransportModuleConstants.RPC_APP_CLIENT_MODULE,
Utils.localHostName(conf),
0,
conf,
diff --git
a/client/src/main/scala/org/apache/celeborn/client/LifecycleManager.scala
b/client/src/main/scala/org/apache/celeborn/client/LifecycleManager.scala
index 8e577a130..a8a9bca86 100644
--- a/client/src/main/scala/org/apache/celeborn/client/LifecycleManager.scala
+++ b/client/src/main/scala/org/apache/celeborn/client/LifecycleManager.scala
@@ -160,7 +160,7 @@ class LifecycleManager(val appUniqueId: String, val conf:
CelebornConf) extends
// init driver celeborn LifecycleManager rpc service
override val rpcEnv: RpcEnv = RpcEnv.create(
RpcNameConstants.LIFECYCLE_MANAGER_SYS,
- TransportModuleConstants.RPC_APP_MODULE,
+ TransportModuleConstants.RPC_LIFECYCLEMANAGER_MODULE,
lifecycleHost,
conf.shuffleManagerPort,
conf,
diff --git
a/common/src/main/java/org/apache/celeborn/common/network/TransportContext.java
b/common/src/main/java/org/apache/celeborn/common/network/TransportContext.java
index aa9543478..625350dc6 100644
---
a/common/src/main/java/org/apache/celeborn/common/network/TransportContext.java
+++
b/common/src/main/java/org/apache/celeborn/common/network/TransportContext.java
@@ -91,6 +91,15 @@ public class TransportContext implements Closeable {
this.channelsLimiter = channelsLimiter;
this.enableHeartbeat = enableHeartbeat;
this.source = source;
+
+ if (null != this.sslFactory) {
+ logger.info(
+ "SSL factory created for module {}, has keys ? {}",
+ conf.getModuleName(),
+ this.sslFactory.hasKeyManagers());
+ } else {
+ logger.info("SSL not enabled for module = {}", conf.getModuleName());
+ }
}
public TransportContext(
diff --git
a/common/src/main/java/org/apache/celeborn/common/network/util/TransportConf.java
b/common/src/main/java/org/apache/celeborn/common/network/util/TransportConf.java
index f0203df4c..c7a0b4497 100644
---
a/common/src/main/java/org/apache/celeborn/common/network/util/TransportConf.java
+++
b/common/src/main/java/org/apache/celeborn/common/network/util/TransportConf.java
@@ -244,8 +244,8 @@ public class TransportConf {
}
public boolean autoSslEnabled() {
- // auto_ssl_enable is supported only for RPC_APP_MODULE
- if (!TransportModuleConstants.RPC_APP_MODULE.equals(module)) return false;
+ // auto_ssl_enable is supported only for RPC_LIFECYCLEMANAGER_MODULE
+ if (!TransportModuleConstants.RPC_LIFECYCLEMANAGER_MODULE.equals(module))
return false;
// auto ssl must be enabled, and there should be no keystore or trust
store configured
boolean autoSslEnabled = celebornConf.isAutoSslEnabled(module);
diff --git
a/common/src/main/java/org/apache/celeborn/common/protocol/TransportModuleConstants.java
b/common/src/main/java/org/apache/celeborn/common/protocol/TransportModuleConstants.java
index 9e813ae3a..6c5bf8fb0 100644
---
a/common/src/main/java/org/apache/celeborn/common/protocol/TransportModuleConstants.java
+++
b/common/src/main/java/org/apache/celeborn/common/protocol/TransportModuleConstants.java
@@ -24,11 +24,21 @@ public class TransportModuleConstants {
// RPC module used by the application components to communicate with each
other
// This is used only at the application side.
+ // This is interally further split into RPC_LIFECYCLEMANAGER_MODULE and
+ // RPC_APP_CLIENT_MODULE - both of which inherit from RPC_APP_MODULE
+ // So for users, there is only RPC_APP_MODULE
public static final String RPC_APP_MODULE = "rpc_app";
// RPC module used to communicate with/between server components
// This is used both at server (master/worker) and application side.
public static final String RPC_SERVICE_MODULE = "rpc_service";
+ // See RPC_APP_MODULE for details - both RPC_LIFECYCLEMANAGER_MODULE and
RPC_APP_CLIENT_MODULE
+ // are internal modules, and transport configs are not expected for these.
+ // For example, auto-ssl requires both RPC_LIFECYCLEMANAGER_MODULE and
RPC_APP_CLIENT_MODULE
+ // to be in sync, though it is enabled only in RPC_LIFECYCLEMANAGER_MODULE
+ public static final String RPC_LIFECYCLEMANAGER_MODULE =
"rpc_app_lifecyclemanager";
+ public static final String RPC_APP_CLIENT_MODULE = "rpc_app_client";
+
// Both RPC_APP and RPC_SERVER fallsback to earlier RPC_MODULE for backward
// compatibility
@Deprecated public static final String RPC_MODULE = "rpc";
diff --git
a/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
b/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
index 9e7d51d97..7c9d6ecb7 100644
--- a/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
+++ b/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala
@@ -87,6 +87,13 @@ class CelebornConf(loadDefaults: Boolean) extends Cloneable
with Logging with Se
this
}
+ private def warnIfInternalTransportModule(module: String, key: String): Unit
= {
+ if (INTERNAL_TRANSPORT_MODULES.contains(module)) {
+ log.warn(s"$key configured for internal transport module $module, " +
+ s"should be using ${INTERNAL_TRANSPORT_MODULES(module)} instead")
+ }
+ }
+
private def getTransportConfContainsImpl[T](module: String, entry:
ConfigEntry[T]): Boolean = {
var currentModule = module
@@ -94,6 +101,7 @@ class CelebornConf(loadDefaults: Boolean) extends Cloneable
with Logging with Se
val key = entry.key.replace("<module>", currentModule)
val opt = getOption(key)
if (opt.isDefined) {
+ warnIfInternalTransportModule(currentModule, key)
return true
}
@@ -120,6 +128,7 @@ class CelebornConf(loadDefaults: Boolean) extends Cloneable
with Logging with Se
val key = configEntry.key.replace("<module>", currentModule)
val opt = getOption(key)
if (opt.isDefined) {
+ warnIfInternalTransportModule(currentModule, key)
return opt.map(converter)
}
@@ -146,6 +155,22 @@ class CelebornConf(loadDefaults: Boolean) extends
Cloneable with Logging with Se
}
}
+ def setTransportConfIfMissing[T](
+ module: String,
+ configEntry: ConfigEntry[T],
+ value: String): Unit = {
+
+ if (getTransportConfContainsImpl(module, configEntry)) {
+ return
+ }
+
+ val key = configEntry.key.replace(
+ "<module>",
+ // if this is an internal module. map it to the exposed module - else
use the same
+ INTERNAL_TRANSPORT_MODULES.getOrElse(module, module))
+ set(key, value)
+ }
+
def set[T](entry: ConfigEntry[T], value: T): CelebornConf = {
set(entry.key, entry.stringConverter(value))
this
@@ -1371,12 +1396,24 @@ class CelebornConf(loadDefaults: Boolean) extends
Cloneable with Logging with Se
object CelebornConf extends Logging {
- val TRANSPORT_MODULE_FALLBACKS = Map(
- "rpc_service" -> "rpc",
- "rpc_app" -> "rpc",
+ val TRANSPORT_MODULE_FALLBACKS: Map[String, String] = Map(
+ TransportModuleConstants.RPC_SERVICE_MODULE ->
TransportModuleConstants.RPC_MODULE,
+ TransportModuleConstants.RPC_APP_MODULE ->
TransportModuleConstants.RPC_MODULE,
+
+ // Internally, RPC_APP_MODULE is split into RPC_LIFECYCLEMANAGER_MODULE and
+ // RPC_APP_CLIENT_MODULE, though this is not exposed to users.
+ TransportModuleConstants.RPC_LIFECYCLEMANAGER_MODULE ->
TransportModuleConstants.RPC_APP_MODULE,
+ TransportModuleConstants.RPC_APP_CLIENT_MODULE ->
TransportModuleConstants.RPC_APP_MODULE,
+
// only for testing
"test_child_module" -> "test_parent_module")
+ // The keys are modules are internal to Celeborn, and users are not expected
to directly
+ // configure them. The values give the user exposed module.
+ val INTERNAL_TRANSPORT_MODULES: Map[String, String] = Map(
+ TransportModuleConstants.RPC_LIFECYCLEMANAGER_MODULE ->
TransportModuleConstants.RPC_APP_MODULE,
+ TransportModuleConstants.RPC_APP_CLIENT_MODULE ->
TransportModuleConstants.RPC_APP_MODULE)
+
/**
* Holds information about keys that have been deprecated and do not have a
replacement.
*
diff --git a/common/src/main/scala/org/apache/celeborn/common/rpc/RpcEnv.scala
b/common/src/main/scala/org/apache/celeborn/common/rpc/RpcEnv.scala
index 90847f71b..0e7573dc0 100644
--- a/common/src/main/scala/org/apache/celeborn/common/rpc/RpcEnv.scala
+++ b/common/src/main/scala/org/apache/celeborn/common/rpc/RpcEnv.scala
@@ -193,7 +193,14 @@ private[celeborn] case class RpcEnvConfig(
port: Int,
numUsableCores: Int,
securityContext: Option[RpcSecurityContext]) {
- assert(TransportModuleConstants.RPC_APP_MODULE == transportModule ||
- TransportModuleConstants.RPC_SERVICE_MODULE == transportModule ||
- TransportModuleConstants.RPC_MODULE == transportModule)
+ assert(RpcEnvConfig.VALID_TRANSPORT_MODULES.contains(transportModule))
+}
+
+object RpcEnvConfig {
+ private val VALID_TRANSPORT_MODULES = Set(
+ TransportModuleConstants.RPC_APP_MODULE,
+ TransportModuleConstants.RPC_MODULE,
+ TransportModuleConstants.RPC_APP_CLIENT_MODULE,
+ TransportModuleConstants.RPC_LIFECYCLEMANAGER_MODULE,
+ TransportModuleConstants.RPC_SERVICE_MODULE)
}
diff --git a/common/src/main/scala/org/apache/celeborn/common/util/Utils.scala
b/common/src/main/scala/org/apache/celeborn/common/util/Utils.scala
index 07e9fe30b..7eef1f4f0 100644
--- a/common/src/main/scala/org/apache/celeborn/common/util/Utils.scala
+++ b/common/src/main/scala/org/apache/celeborn/common/util/Utils.scala
@@ -494,10 +494,19 @@ object Utils extends Logging {
// assuming we have all the machine's cores).
// NB: Only set if serverThreads/clientThreads not already set.
val numThreads = defaultNumThreads(numUsableCores)
- conf.setIfMissing(s"celeborn.$module.io.serverThreads",
numThreads.toString)
- conf.setIfMissing(s"celeborn.$module.io.clientThreads",
numThreads.toString)
+ conf.setTransportConfIfMissing(
+ module,
+ CelebornConf.NETWORK_IO_SERVER_THREADS,
+ numThreads.toString)
+ conf.setTransportConfIfMissing(
+ module,
+ CelebornConf.NETWORK_IO_CLIENT_THREADS,
+ numThreads.toString)
if (TransportModuleConstants.PUSH_MODULE == module) {
- conf.setIfMissing(s"celeborn.$module.io.numConnectionsPerPeer",
numThreads.toString)
+ conf.setTransportConfIfMissing(
+ module,
+ CelebornConf.NETWORK_IO_NUM_CONNECTIONS_PER_PEER,
+ numThreads.toString)
}
new TransportConf(module, conf)
diff --git
a/common/src/test/java/org/apache/celeborn/common/network/AutoSSLRpcIntegrationSuite.java
b/common/src/test/java/org/apache/celeborn/common/network/AutoSSLRpcIntegrationSuite.java
index 7e6af4548..a975ed069 100644
---
a/common/src/test/java/org/apache/celeborn/common/network/AutoSSLRpcIntegrationSuite.java
+++
b/common/src/test/java/org/apache/celeborn/common/network/AutoSSLRpcIntegrationSuite.java
@@ -30,12 +30,16 @@ import
org.apache.celeborn.common.protocol.TransportModuleConstants;
public class AutoSSLRpcIntegrationSuite extends RpcIntegrationSuiteJ {
@BeforeClass
public static void setUp() throws Exception {
- // change it to client module
- RpcIntegrationSuiteJ.TEST_MODULE = TransportModuleConstants.RPC_APP_MODULE;
+ // change it to lifecycle manager module
+ RpcIntegrationSuiteJ.TEST_MODULE =
TransportModuleConstants.RPC_LIFECYCLEMANAGER_MODULE;
// set up SSL for TEST_MODULE
RpcIntegrationSuiteJ.initialize(
TestHelper.updateCelebornConfWithMap(
- new CelebornConf(),
SslSampleConfigs.createAutoSslConfigForModule(TEST_MODULE)));
+ new CelebornConf(),
+ // to validate, we are using
TransportModuleConstants.RPC_APP_MODULE, to mirror how
+ // users will configure
+ SslSampleConfigs.createAutoSslConfigForModule(
+ TransportModuleConstants.RPC_APP_MODULE)));
}
@AfterClass
diff --git
a/common/src/test/java/org/apache/celeborn/common/network/ssl/SslConnectivitySuiteJ.java
b/common/src/test/java/org/apache/celeborn/common/network/ssl/SslConnectivitySuiteJ.java
index 88b9f293c..52fadb69c 100644
---
a/common/src/test/java/org/apache/celeborn/common/network/ssl/SslConnectivitySuiteJ.java
+++
b/common/src/test/java/org/apache/celeborn/common/network/ssl/SslConnectivitySuiteJ.java
@@ -299,8 +299,6 @@ public class SslConnectivitySuiteJ {
return conf;
};
- // checking nettyssl == false at both client and server does not make
sense in this context
- // it is the same cert for both :)
testSuccessfulConnectivity(true, true, true, updateConf, updateConf);
testSuccessfulConnectivity(true, true, false, updateConf, updateConf);
testSuccessfulConnectivity(true, false, true, updateConf, updateConf);
@@ -316,6 +314,7 @@ public class SslConnectivitySuiteJ {
throws Exception {
testSuccessfulConnectivity(
+ TEST_MODULE,
TEST_MODULE,
enableSsl,
primaryConfigForServer,
@@ -325,7 +324,8 @@ public class SslConnectivitySuiteJ {
}
private void testSuccessfulConnectivity(
- String module,
+ String serverModule,
+ String clientModule,
boolean enableSsl,
boolean primaryConfigForServer,
boolean primaryConfigForClient,
@@ -336,9 +336,9 @@ public class SslConnectivitySuiteJ {
new TestTransportState(
DEFAULT_HANDLER,
createTransportConf(
- module, enableSsl, primaryConfigForServer,
postProcessServerConf),
+ serverModule, enableSsl, primaryConfigForServer,
postProcessServerConf),
createTransportConf(
- module, enableSsl, primaryConfigForClient,
postProcessClientConf));
+ clientModule, enableSsl, primaryConfigForClient,
postProcessClientConf));
TransportClient client = state.createClient()) {
String msg = " hi ";
@@ -396,11 +396,15 @@ public class SslConnectivitySuiteJ {
@Test
public void testAutoSslConnectivity() throws Exception {
- String module = TransportModuleConstants.RPC_MODULE;
- final Function<CelebornConf, CelebornConf> updateConf =
+ // Mirror how user will configure - so configure module based on
RPC_APP_MODULE
+ // while create the server transport config for lifecycle manager (to
match driver),
+ // and client to app_client similar to executors
+
+ final Function<CelebornConf, CelebornConf> updateServerConf =
conf -> {
// return a new config
+ String module = TransportModuleConstants.RPC_APP_MODULE;
CelebornConf celebornConf = new CelebornConf();
celebornConf.set("celeborn.ssl." + module + ".enabled", "true");
celebornConf.set("celeborn.ssl." + module + ".protocol", "TLSv1.2");
@@ -408,9 +412,23 @@ public class SslConnectivitySuiteJ {
return celebornConf;
};
- // checking nettyssl == false at both client and server does not make
sense in this context
- // it is the same cert for both :)
+ final Function<CelebornConf, CelebornConf> updateClientConf =
+ conf -> {
+ // return a new config
+ String module = TransportModuleConstants.RPC_APP_MODULE;
+ CelebornConf celebornConf = new CelebornConf();
+ celebornConf.set("celeborn.ssl." + module + ".enabled", "true");
+ celebornConf.set("celeborn.ssl." + module + ".protocol", "TLSv1.2");
+ return celebornConf;
+ };
+
testSuccessfulConnectivity(
- TransportModuleConstants.RPC_APP_MODULE, true, true, true, updateConf,
updateConf);
+ TransportModuleConstants.RPC_LIFECYCLEMANAGER_MODULE,
+ TransportModuleConstants.RPC_APP_CLIENT_MODULE,
+ true,
+ true,
+ true,
+ updateServerConf,
+ updateClientConf);
}
}
diff --git
a/common/src/test/java/org/apache/celeborn/common/network/util/TransportConfSuiteJ.java
b/common/src/test/java/org/apache/celeborn/common/network/util/TransportConfSuiteJ.java
index 2f4b1e0ad..e03d1b536 100644
---
a/common/src/test/java/org/apache/celeborn/common/network/util/TransportConfSuiteJ.java
+++
b/common/src/test/java/org/apache/celeborn/common/network/util/TransportConfSuiteJ.java
@@ -114,7 +114,7 @@ public class TransportConfSuiteJ {
// for rpc_app module, it can be enabled.
TransportConf conf =
new TransportConf(
- TransportModuleConstants.RPC_APP_MODULE,
+ TransportModuleConstants.RPC_LIFECYCLEMANAGER_MODULE,
TestHelper.updateCelebornConfWithMap(
new CelebornConf(),
SslSampleConfigs.createAutoSslConfigForModule(
diff --git
a/common/src/test/scala/org/apache/celeborn/common/rpc/RpcEnvSuite.scala
b/common/src/test/scala/org/apache/celeborn/common/rpc/RpcEnvSuite.scala
index 4843cf744..d1ab75f37 100644
--- a/common/src/test/scala/org/apache/celeborn/common/rpc/RpcEnvSuite.scala
+++ b/common/src/test/scala/org/apache/celeborn/common/rpc/RpcEnvSuite.scala
@@ -666,7 +666,9 @@ abstract class RpcEnvSuite extends CelebornFunSuite {
test("port conflict") {
val anotherEnv = createRpcEnv(createCelebornConf(), "remote",
env.address.port)
try {
- assert(anotherEnv.address.port != env.address.port)
+ assert(
+ anotherEnv.address.port != env.address.port,
+ s"new port = ${anotherEnv.address.port}, env port =
${env.address.port}")
} finally {
anotherEnv.shutdown()
anotherEnv.awaitTermination()