This is an automated email from the ASF dual-hosted git repository.

chengpan pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/incubator-celeborn.git


The following commit(s) were added to refs/heads/main by this push:
     new 374d735ae [CELEBORN-724] Fix the compatibility of 
HeartbeatFromApplicationRespo…
374d735ae is described below

commit 374d735ae598240aa8722d09dbfc80f11c3b30e5
Author: zhongqiang.czq <[email protected]>
AuthorDate: Wed Jun 28 16:04:18 2023 +0800

    [CELEBORN-724] Fix the compatibility of HeartbeatFromApplicationRespo…
    
    …nse with lower versions
    
    ### What changes were proposed in this pull request?
    The master side will check HeartbeatFromApplication's reply field. if reply 
is true then it replies HeartbeatFromApplicationResponse otherwise 
OneWayMessageResponse.
    
    The reply field is default false before the version 0.2.1, so master can be 
compatible with older client version
    
    ### Why are the changes needed?
    Before the version `0.2.1`, the response of HeartbeatFromApplication is` 
OneWayMessageResponse`, but from `0.3.0`, the response of 
HeartbeatFromApplication is modified to `HeartbeatFromApplicationResponse`.
    if the version of `client side `is `0.2.1` and the version of `server side 
is 0.3.0`, the `compatiblity issue `will occur.
    The following compatiblity error will be printted.
    
    ``` java
    java.io.InvalidObjectException: enum constant 
HEARTBEAT_FROM_APPLICATION_RESPONSE does not exist in class 
org.apache.celeborn.common.protocol.MessageType
            at java.io.ObjectInputStream.readEnum(ObjectInputStream.java:2157) 
~[?:1.8.0_362]
            at 
java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1662) 
~[?:1.8.0_362]
            at 
java.io.ObjectInputStream.defaultReadFields(ObjectInputStream.java:2430) 
~[?:1.8.0_362]
            at 
java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2354) 
~[?:1.8.0_362]
            at 
java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2212) 
~[?:1.8.0_362]
            at 
java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1668) 
~[?:1.8.0_362]
            at java.io.ObjectInputStream.readObject(ObjectInputStream.java:502) 
~[?:1.8.0_362]
            at java.io.ObjectInputStream.readObject(ObjectInputStream.java:460) 
~[?:1.8.0_362]
            at 
org.apache.celeborn.common.serializer.JavaDeserializationStream.readObject(JavaSerializer.scala:76)
 ~[celeborn-client-spark-3-shaded_2.12-0.2.1-incubating.jar:?]
    ```
    ``` java
    Caused by: java.lang.ClassCastException: Cannot cast 
org.apache.celeborn.common.protocol.message.ControlMessages$HeartbeatFromApplicationResponse
 to 
org.apache.celeborn.common.protocol.message.ControlMessages$OneWayMessageResponse$
            at java.lang.Class.cast(Class.java:3369) ~[?:1.8.0_362]
            at scala.concurrent.Future.$anonfun$mapTo$1(Future.scala:500) 
~[scala-library-2.12.15.jar:?]
            at scala.util.Success.$anonfun$map$1(Try.scala:255) 
~[scala-library-2.12.15.jar:?]
            at scala.util.Success.map(Try.scala:213) 
~[scala-library-2.12.15.jar:?]
            at scala.concurrent.Future.$anonfun$map$1(Future.scala:292) 
~[scala-library-2.12.15.jar:?]
            at scala.concurrent.impl.Promise.liftedTree1$1(Promise.scala:33) 
~[scala-library-2.12.15.jar:?]
            at 
scala.concurrent.impl.Promise.$anonfun$transform$1(Promise.scala:33) 
~[scala-library-2.12.15.jar:?]
            at scala.concurrent.impl.CallbackRunnable.run(Promise.scala:64) 
~[scala-library-2.12.15.jar:?]
            at 
scala.concurrent.BatchingExecutor$Batch.processBatch$1(BatchingExecutor.scala:67)
 ~[scala-library-2.12.15.jar:?]
            at 
scala.concurrent.BatchingExecutor$Batch.$anonfun$run$1(BatchingExecutor.scala:82)
 ~[scala-library-2.12.15.jar:?]
            at 
scala.runtime.java8.JFunction0$mcV$sp.apply(JFunction0$mcV$sp.java:23) 
~[scala-library-2.12.15.jar:?]
            at 
scala.concurrent.BlockContext$.withBlockContext(BlockContext.scala:85) 
~[scala-library-2.12.15.jar:?]
            at 
scala.concurrent.BatchingExecutor$Batch.run(BatchingExecutor.scala:59) 
~[scala-library-2.12.15.jar:?]
            at 
scala.concurrent.Future$InternalCallbackExecutor$.unbatchedExecute(Future.scala:875)
 ~[scala-library-2.12.15.jar:?]
            at 
scala.concurrent.BatchingExecutor.execute(BatchingExecutor.scala:110) 
~[scala-library-2.12.15.jar:?]
            at 
scala.concurrent.BatchingExecutor.execute$(BatchingExecutor.scala:107) 
~[scala-library-2.12.15.jar:?]
            at 
scala.concurrent.Future$InternalCallbackExecutor$.execute(Future.scala:873) 
~[scala-library-2.12.15.jar:?]
            at 
scala.concurrent.impl.CallbackRunnable.executeWithValue(Promise.scala:72) 
~[scala-library-2.12.15.jar:?]
            at 
scala.concurrent.impl.Promise$DefaultPromise.$anonfun$tryComplete$1(Promise.scala:288)
 ~[scala-library-2.12.15.jar:?]
            at 
scala.concurrent.impl.Promise$DefaultPromise.$anonfun$tryComplete$1$adapted(Promise.scala:288)
 ~[scala-library-2.12.15.jar:?]
            at 
scala.concurrent.impl.Promise$DefaultPromise.tryComplete(Promise.scala:288) 
~[scala-library-2.12.15.jar:?]
            at scala.concurrent.Promise.trySuccess(Promise.scala:94) 
~[scala-library-2.12.15.jar:?]
            at scala.concurrent.Promise.trySuccess$(Promise.scala:94) 
~[scala-library-2.12.15.jar:?]
            at 
scala.concurrent.impl.Promise$DefaultPromise.trySuccess(Promise.scala:187) 
~[scala-library-2.12.15.jar:?]
            at 
org.apache.celeborn.common.rpc.netty.NettyRpcEnv.onSuccess$1(NettyRpcEnv.scala:218)
 ~[celeborn-client-spark-3-shaded_2.12-0.2.1-incubating.jar:?]
    ```
    
    ### Does this PR introduce _any_ user-facing change?
    
    No
    
    ### How was this patch tested?
    The pr is tested manually and the testing process is as follows:
    1. server side is deploy using the code of latest branch-0.3.
    2. spark client is deploy the version of 0.2.1, then run spark-sql to 
execute  3 tpcds queries( query1.sql/querey2/quere3.sql whose datasize is 1T), 
finnally verify that the queries are executed successfully and no above 
compatiblity error printted
    3. spark client is deploy the version of 0.3.0,  then run spark-sql to 
execute 3 tpcds queries( query1.sql/querey2/quere3.sql whose datasize is 1T), 
finnally verify that the queries are executed successfully and no above 
compatiblity error printted
    
    This patch had conflicts when merged, resolved by
    Committer: Cheng Pan <[email protected]>
    
    Closes #1635 from zhongqiangczq/heartbeat2.
    
    Authored-by: zhongqiang.czq <[email protected]>
    Signed-off-by: Cheng Pan <[email protected]>
---
 .../celeborn/client/ApplicationHeartbeater.scala   |  3 ++-
 common/src/main/proto/TransportMessages.proto      |  1 +
 .../common/protocol/message/ControlMessages.scala  | 10 +++++---
 docs/migration.md                                  | 26 ++++----------------
 .../celeborn/service/deploy/master/Master.scala    | 28 +++++++++++++++-------
 5 files changed, 35 insertions(+), 33 deletions(-)

diff --git 
a/client/src/main/scala/org/apache/celeborn/client/ApplicationHeartbeater.scala 
b/client/src/main/scala/org/apache/celeborn/client/ApplicationHeartbeater.scala
index a3465cbd1..dd8a5d389 100644
--- 
a/client/src/main/scala/org/apache/celeborn/client/ApplicationHeartbeater.scala
+++ 
b/client/src/main/scala/org/apache/celeborn/client/ApplicationHeartbeater.scala
@@ -57,7 +57,8 @@ class ApplicationHeartbeater(
                 tmpTotalWritten,
                 tmpTotalFileCount,
                 workerStatusTracker.getNeedCheckedWorkers().toList.asJava,
-                ZERO_UUID)
+                ZERO_UUID,
+                true)
             val response = requestHeartbeat(appHeartbeat)
             if (response.statusCode == StatusCode.SUCCESS) {
               logDebug("Successfully send app heartbeat.")
diff --git a/common/src/main/proto/TransportMessages.proto 
b/common/src/main/proto/TransportMessages.proto
index c3ebb92a5..b4f6a0c60 100644
--- a/common/src/main/proto/TransportMessages.proto
+++ b/common/src/main/proto/TransportMessages.proto
@@ -290,6 +290,7 @@ message PbHeartbeatFromApplication {
   int64 fileCount = 3 ;
   string requestId = 4;
   repeated PbWorkerInfo needCheckedWorkerList = 5;
+  bool shouldResponse = 6;
 }
 
 message PbHeartbeatFromApplicationResponse {
diff --git 
a/common/src/main/scala/org/apache/celeborn/common/protocol/message/ControlMessages.scala
 
b/common/src/main/scala/org/apache/celeborn/common/protocol/message/ControlMessages.scala
index ef42e6603..9354bee62 100644
--- 
a/common/src/main/scala/org/apache/celeborn/common/protocol/message/ControlMessages.scala
+++ 
b/common/src/main/scala/org/apache/celeborn/common/protocol/message/ControlMessages.scala
@@ -323,7 +323,8 @@ object ControlMessages extends Logging {
       totalWritten: Long,
       fileCount: Long,
       needCheckedWorkerList: util.List[WorkerInfo],
-      override var requestId: String = ZERO_UUID) extends MasterRequestMessage
+      override var requestId: String = ZERO_UUID,
+      shouldResponse: Boolean = false) extends MasterRequestMessage
 
   case class HeartbeatFromApplicationResponse(
       statusCode: StatusCode,
@@ -619,7 +620,8 @@ object ControlMessages extends Logging {
           totalWritten,
           fileCount,
           needCheckedWorkerList,
-          requestId) =>
+          requestId,
+          shouldResponse) =>
       val payload = PbHeartbeatFromApplication.newBuilder()
         .setAppId(appId)
         .setRequestId(requestId)
@@ -627,6 +629,7 @@ object ControlMessages extends Logging {
         .setFileCount(fileCount)
         .addAllNeedCheckedWorkerList(needCheckedWorkerList.asScala.map(
           PbSerDeUtils.toPbWorkerInfo(_, true)).toList.asJava)
+        .setShouldResponse(shouldResponse)
         .build().toByteArray
       new TransportMessage(MessageType.HEARTBEAT_FROM_APPLICATION, payload)
 
@@ -972,7 +975,8 @@ object ControlMessages extends Logging {
           new util.ArrayList[WorkerInfo](
             pbHeartbeatFromApplication.getNeedCheckedWorkerListList.asScala
               .map(PbSerDeUtils.fromPbWorkerInfo).toList.asJava),
-          pbHeartbeatFromApplication.getRequestId)
+          pbHeartbeatFromApplication.getRequestId,
+          pbHeartbeatFromApplication.getShouldResponse)
 
       case HEARTBEAT_FROM_APPLICATION_RESPONSE =>
         val pbHeartbeatFromApplicationResponse =
diff --git a/docs/migration.md b/docs/migration.md
index d19647a0a..727817a32 100644
--- a/docs/migration.md
+++ b/docs/migration.md
@@ -19,7 +19,10 @@ license: |
 
 # Migration Guide
 
-## Upgrading from 0.2.1 to 0.3.0
+## Upgrading from 0.2 to 0.3
+
+ - Celeborn 0.2 Client is compatible with 0.3 Master/Server, it allows to 
upgrade Master/Worker first then Client.
+   Note that: It's strongly recommended to use the same version of Client and 
Celeborn Master/Worker in production.
 
  - From 0.3.0 on the default value for 
`celeborn.client.push.replicate.enabled` is changed from `true` to `false`, 
users
    who want replication on should explicitly enable replication. For example, 
to enable replication for Spark
@@ -37,25 +40,6 @@ license: |
  - Since 0.3.0, the Celeborn Master URL schema is changed from `rss://` to 
`celeborn://`, for users who start Worker by
    `sbin/start-worker.sh rss://<master-host>:<master-port>`, should migrate to 
`sbin/start-worker.sh celeborn://<master-host>:<master-port>`.
 
- - When using 0.2.1 as client side and 0.3.0 as server side, you may see the 
following Exception in LifecycleManger's
-   log. You can safely ignore the log, it's caused by the behavior change when 
Master receives heartbeat from Application.
-
-    ??? warning "logs"
-        ```
-        23/06/20 18:12:30 WARN TransportChannelHandler: Exception in 
connection from /192.168.1.16:9097
-        java.io.InvalidObjectException: enum constant 
HEARTBEAT_FROM_APPLICATION_RESPONSE does not exist in class 
org.apache.celeborn.common.protocol.MessageType
-            at java.io.ObjectInputStream.readEnum(ObjectInputStream.java:2157)
-            at 
java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1662)
-            at 
java.io.ObjectInputStream.defaultReadFields(ObjectInputStream.java:2430)
-            at 
java.io.ObjectInputStream.readSerialData(ObjectInputStream.java:2354)
-            at 
java.io.ObjectInputStream.readOrdinaryObject(ObjectInputStream.java:2212)
-            at 
java.io.ObjectInputStream.readObject0(ObjectInputStream.java:1668)
-            at java.io.ObjectInputStream.readObject(ObjectInputStream.java:502)
-            at java.io.ObjectInputStream.readObject(ObjectInputStream.java:460)
-            at 
org.apache.celeborn.common.serializer.JavaDeserializationStream.readObject(JavaSerializer.scala:76)
-            at 
org.apache.celeborn.common.serializer.JavaSerializerInstance.deserialize(JavaSerializer.scala:110)
-        ```
-
  - Since 0.3.0, Celeborn supports overriding Hadoop 
configuration(`core-site.xml`, `hdfs-site.xml`, etc.) from Celeborn 
configuration with the additional prefix `celeborn.hadoop.`. 
    On Spark client side, user should set Hadoop configuration like 
`spark.celeborn.hadoop.foo=bar`, note that `spark.hadoop.foo=bar` does not take 
effect;
-   on Flink client and Celeborn Master/Worker side, user should set like 
`celeborn.hadoop.foo=bar`.
\ No newline at end of file
+   on Flink client and Celeborn Master/Worker side, user should set like 
`celeborn.hadoop.foo=bar`.
diff --git 
a/master/src/main/scala/org/apache/celeborn/service/deploy/master/Master.scala 
b/master/src/main/scala/org/apache/celeborn/service/deploy/master/Master.scala
index 73182d50a..e001a6f1c 100644
--- 
a/master/src/main/scala/org/apache/celeborn/service/deploy/master/Master.scala
+++ 
b/master/src/main/scala/org/apache/celeborn/service/deploy/master/Master.scala
@@ -230,7 +230,13 @@ private[celeborn] class Master(
   }
 
   override def receiveAndReply(context: RpcCallContext): PartialFunction[Any, 
Unit] = {
-    case HeartbeatFromApplication(appId, totalWritten, fileCount, 
localBlacklist, requestId) =>
+    case HeartbeatFromApplication(
+          appId,
+          totalWritten,
+          fileCount,
+          localBlacklist,
+          requestId,
+          shouldResponse) =>
       logDebug(s"Received heartbeat from app $appId")
       executeWithLeaderChecker(
         context,
@@ -240,7 +246,8 @@ private[celeborn] class Master(
           totalWritten,
           fileCount,
           localBlacklist,
-          requestId))
+          requestId,
+          shouldResponse))
 
     case pbRegisterWorker: PbRegisterWorker =>
       val requestId = pbRegisterWorker.getRequestId
@@ -665,7 +672,8 @@ private[celeborn] class Master(
       totalWritten: Long,
       fileCount: Long,
       needCheckedWorkerList: util.List[WorkerInfo],
-      requestId: String): Unit = {
+      requestId: String,
+      shouldResponse: Boolean): Unit = {
     statusSystem.handleAppHeartbeat(
       appId,
       totalWritten,
@@ -674,11 +682,15 @@ private[celeborn] class Master(
       requestId)
     // unknown workers will retain in needCheckedWorkerList
     needCheckedWorkerList.removeAll(workersSnapShot)
-    context.reply(HeartbeatFromApplicationResponse(
-      StatusCode.SUCCESS,
-      new util.ArrayList(statusSystem.blacklist),
-      needCheckedWorkerList,
-      shutdownWorkerSnapshot))
+    if (shouldResponse) {
+      context.reply(HeartbeatFromApplicationResponse(
+        StatusCode.SUCCESS,
+        new util.ArrayList(statusSystem.blacklist),
+        needCheckedWorkerList,
+        shutdownWorkerSnapshot))
+    } else {
+      context.reply(OneWayMessageResponse)
+    }
   }
 
   private def computeUserResourceConsumption(userIdentifier: UserIdentifier)

Reply via email to