This is an automated email from the ASF dual-hosted git repository.
rexxiong pushed a commit to branch branch-0.5
in repository https://gitbox.apache.org/repos/asf/celeborn.git
The following commit(s) were added to refs/heads/branch-0.5 by this push:
new 41387aa06 [CELEBORN-1337] Remove unused fields from
HeartbeatFromApplicationResponse
41387aa06 is described below
commit 41387aa06e1746b3c34fbd1ec8eabf0a82fc808c
Author: Erik.fang <[email protected]>
AuthorDate: Thu Jun 13 10:15:00 2024 +0800
[CELEBORN-1337] Remove unused fields from HeartbeatFromApplicationResponse
as discussed in https://github.com/apache/celeborn/pull/2398, this PR
removed unused fields from HeartbeatFromApplicationResponse, without adding
WorkerId Type
Closes #2529 from
ErikFang/remove-unused-fields-HeartbeatFromApplicationResponse.
Authored-by: Erik.fang <[email protected]>
Signed-off-by: Shuang <[email protected]>
(cherry picked from commit 5323c1d009f1c17135eff65179b93c7f06326174)
Signed-off-by: Shuang <[email protected]>
---
.../celeborn/client/ApplicationHeartbeater.scala | 2 ++
.../apache/celeborn/client/LifecycleManager.scala | 2 ++
.../common/protocol/message/ControlMessages.scala | 21 ++++++++++----------
.../apache/celeborn/common/util/PbSerDeUtils.scala | 23 ++++++++++++----------
.../celeborn/common/util/PbSerDeUtilsTest.scala | 4 ++--
.../celeborn/service/deploy/master/Master.scala | 2 ++
6 files changed, 32 insertions(+), 22 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 5101d61f1..d3af38183 100644
---
a/client/src/main/scala/org/apache/celeborn/client/ApplicationHeartbeater.scala
+++
b/client/src/main/scala/org/apache/celeborn/client/ApplicationHeartbeater.scala
@@ -54,6 +54,8 @@ class ApplicationHeartbeater(
val (tmpTotalWritten, tmpTotalFileCount) = shuffleMetrics()
logInfo("Send app heartbeat with " +
s"written: ${Utils.bytesToString(tmpTotalWritten)}, file count:
$tmpTotalFileCount")
+ // UserResourceConsumption and DiskInfo are eliminated from
WorkerInfo
+ // during serialization of HeartbeatFromApplication
val appHeartbeat =
HeartbeatFromApplication(
appId,
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 e3745b431..d55aee47d 100644
--- a/client/src/main/scala/org/apache/celeborn/client/LifecycleManager.scala
+++ b/client/src/main/scala/org/apache/celeborn/client/LifecycleManager.scala
@@ -1529,6 +1529,8 @@ class LifecycleManager(val appUniqueId: String, val conf:
CelebornConf) extends
} else {
Set.empty[WorkerInfo]
}
+ // UserResourceConsumption and DiskInfo are eliminated from WorkerInfo
+ // during serialization of RequestSlots
val req =
RequestSlots(
appUniqueId,
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 9491ba29d..4330d8a23 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
@@ -291,10 +291,10 @@ object ControlMessages extends Logging {
workersToRemove: util.List[WorkerInfo],
requestId: String): PbWorkerExclude = PbWorkerExclude.newBuilder()
.addAllWorkersToAdd(workersToAdd.asScala.map { workerInfo =>
- PbSerDeUtils.toPbWorkerInfo(workerInfo, true)
+ PbSerDeUtils.toPbWorkerInfo(workerInfo, true, false)
}.toList.asJava)
.addAllWorkersToRemove(workersToRemove.asScala.map { workerInfo =>
- PbSerDeUtils.toPbWorkerInfo(workerInfo, true)
+ PbSerDeUtils.toPbWorkerInfo(workerInfo, true, false)
}.toList.asJava)
.setRequestId(requestId)
.build()
@@ -403,7 +403,7 @@ object ControlMessages extends Logging {
PbRemoveWorkersUnavailableInfo.newBuilder()
.setRequestId(requestId)
.addAllWorkerInfo(unavailable.asScala.map { workerInfo =>
- PbSerDeUtils.toPbWorkerInfo(workerInfo, true)
+ PbSerDeUtils.toPbWorkerInfo(workerInfo, true, false)
}.toList.asJava)
.build()
}
@@ -417,7 +417,7 @@ object ControlMessages extends Logging {
.setRequestId(requestId)
.setWorkerEventType(WorkerEventType.valueOf(eventType))
.addAllWorkers(workers.asScala.map { workerInfo =>
- PbSerDeUtils.toPbWorkerInfo(workerInfo, true)
+ PbSerDeUtils.toPbWorkerInfo(workerInfo, true, false)
}.toList.asJava)
.build()
}
@@ -599,7 +599,8 @@ object ControlMessages extends Logging {
.setRequestId(requestId)
.setAvailableStorageTypes(availableStorageTypes)
.setUserIdentifier(PbSerDeUtils.toPbUserIdentifier(userIdentifier))
-
.addAllExcludedWorkerSet(excludedWorkerSet.map(PbSerDeUtils.toPbWorkerInfo(_,
true)).asJava)
+ .addAllExcludedWorkerSet(excludedWorkerSet.map(
+ PbSerDeUtils.toPbWorkerInfo(_, true, true)).asJava)
.setPacked(packed)
.build().toByteArray
new TransportMessage(MessageType.REQUEST_SLOTS, payload)
@@ -733,7 +734,7 @@ object ControlMessages extends Logging {
.setTotalWritten(totalWritten)
.setFileCount(fileCount)
.addAllNeedCheckedWorkerList(needCheckedWorkerList.asScala.map(
- PbSerDeUtils.toPbWorkerInfo(_, true)).toList.asJava)
+ PbSerDeUtils.toPbWorkerInfo(_, true, true)).toList.asJava)
.setShouldResponse(shouldResponse)
.build().toByteArray
new TransportMessage(MessageType.HEARTBEAT_FROM_APPLICATION, payload)
@@ -746,11 +747,11 @@ object ControlMessages extends Logging {
val payload = PbHeartbeatFromApplicationResponse.newBuilder()
.setStatus(statusCode.getValue)
.addAllExcludedWorkers(
- excludedWorkers.asScala.map(PbSerDeUtils.toPbWorkerInfo(_,
true)).toList.asJava)
+ excludedWorkers.asScala.map(PbSerDeUtils.toPbWorkerInfo(_, true,
true)).toList.asJava)
.addAllUnknownWorkers(
- unknownWorkers.asScala.map(PbSerDeUtils.toPbWorkerInfo(_,
true)).toList.asJava)
+ unknownWorkers.asScala.map(PbSerDeUtils.toPbWorkerInfo(_, true,
true)).toList.asJava)
.addAllShuttingWorkers(
- shuttingWorkers.asScala.map(PbSerDeUtils.toPbWorkerInfo(_,
true)).toList.asJava)
+ shuttingWorkers.asScala.map(PbSerDeUtils.toPbWorkerInfo(_, true,
true)).toList.asJava)
.build().toByteArray
new TransportMessage(MessageType.HEARTBEAT_FROM_APPLICATION_RESPONSE,
payload)
@@ -771,7 +772,7 @@ object ControlMessages extends Logging {
case ReportWorkerUnavailable(failed, requestId) =>
val payload = PbReportWorkerUnavailable.newBuilder()
.addAllUnavailable(failed.asScala.map { workerInfo =>
- PbSerDeUtils.toPbWorkerInfo(workerInfo, true)
+ PbSerDeUtils.toPbWorkerInfo(workerInfo, true, false)
}
.toList.asJava)
.setRequestId(requestId).build().toByteArray
diff --git
a/common/src/main/scala/org/apache/celeborn/common/util/PbSerDeUtils.scala
b/common/src/main/scala/org/apache/celeborn/common/util/PbSerDeUtils.scala
index cffe76857..0f83d02fc 100644
--- a/common/src/main/scala/org/apache/celeborn/common/util/PbSerDeUtils.scala
+++ b/common/src/main/scala/org/apache/celeborn/common/util/PbSerDeUtils.scala
@@ -240,10 +240,8 @@ object PbSerDeUtils {
def toPbWorkerInfo(
workerInfo: WorkerInfo,
- eliminateUserResourceConsumption: Boolean): PbWorkerInfo = {
- val diskInfos = workerInfo.diskInfos.values
- val pbDiskInfos = new util.ArrayList[PbDiskInfo]()
- diskInfos.asScala.foreach(diskInfo =>
pbDiskInfos.add(PbSerDeUtils.toPbDiskInfo(diskInfo)))
+ eliminateUserResourceConsumption: Boolean,
+ eliminateDiskInfo: Boolean): PbWorkerInfo = {
val builder = PbWorkerInfo.newBuilder
.setHost(workerInfo.host)
.setRpcPort(workerInfo.rpcPort)
@@ -251,11 +249,16 @@ object PbSerDeUtils {
.setPushPort(workerInfo.pushPort)
.setReplicatePort(workerInfo.replicatePort)
.setInternalPort(workerInfo.internalPort)
- .addAllDisks(pbDiskInfos)
if (!eliminateUserResourceConsumption) {
builder.putAllUserResourceConsumption(
PbSerDeUtils.toPbUserResourceConsumption(workerInfo.userResourceConsumption))
}
+ if (!eliminateDiskInfo) {
+ val diskInfos = workerInfo.diskInfos.values
+ val pbDiskInfos = new util.ArrayList[PbDiskInfo]()
+ diskInfos.asScala.foreach(diskInfo =>
pbDiskInfos.add(PbSerDeUtils.toPbDiskInfo(diskInfo)))
+ builder.addAllDisks(pbDiskInfos)
+ }
builder.build
}
@@ -426,12 +429,12 @@ object PbSerDeUtils {
.setEstimatedPartitionSize(estimatedPartitionSize)
.addAllRegisteredShuffle(registeredShuffle)
.addAllHostnameSet(hostnameSet)
- .addAllExcludedWorkers(excludedWorkers.asScala.map(toPbWorkerInfo(_,
true)).asJava)
+ .addAllExcludedWorkers(excludedWorkers.asScala.map(toPbWorkerInfo(_,
true, false)).asJava)
.addAllManuallyExcludedWorkers(manuallyExcludedWorkers.asScala
- .map(toPbWorkerInfo(_, true)).asJava)
- .addAllWorkerLostEvents(workerLostEvent.asScala.map(toPbWorkerInfo(_,
true)).asJava)
+ .map(toPbWorkerInfo(_, true, false)).asJava)
+ .addAllWorkerLostEvents(workerLostEvent.asScala.map(toPbWorkerInfo(_,
true, false)).asJava)
.putAllAppHeartbeatTime(appHeartbeatTime)
- .addAllWorkers(workers.asScala.map(toPbWorkerInfo(_, true)).asJava)
+ .addAllWorkers(workers.asScala.map(toPbWorkerInfo(_, true,
false)).asJava)
.setPartitionTotalWritten(partitionTotalWritten)
.setPartitionTotalFileCount(partitionTotalFileCount)
// appDiskUsageMetricSnapshots can have null values,
@@ -441,7 +444,7 @@ object PbSerDeUtils {
.putAllLostWorkers(lostWorkers.asScala.map {
case (worker: WorkerInfo, time: java.lang.Long) =>
(worker.toUniqueId(), time)
}.asJava)
- .addAllShutdownWorkers(shutdownWorkers.asScala.map(toPbWorkerInfo(_,
true)).asJava)
+ .addAllShutdownWorkers(shutdownWorkers.asScala.map(toPbWorkerInfo(_,
true, false)).asJava)
.putAllWorkerEventInfos(workerEventInfos.asScala.map {
case (worker, workerEventInfo) =>
(worker.toUniqueId(),
PbSerDeUtils.toPbWorkerEventInfo(workerEventInfo))
diff --git
a/common/src/test/scala/org/apache/celeborn/common/util/PbSerDeUtilsTest.scala
b/common/src/test/scala/org/apache/celeborn/common/util/PbSerDeUtilsTest.scala
index a9ad63e27..5f5247273 100644
---
a/common/src/test/scala/org/apache/celeborn/common/util/PbSerDeUtilsTest.scala
+++
b/common/src/test/scala/org/apache/celeborn/common/util/PbSerDeUtilsTest.scala
@@ -239,8 +239,8 @@ class PbSerDeUtilsTest extends CelebornFunSuite {
}
test("fromAndToPbWorkerInfo") {
- val pbWorkerInfo = PbSerDeUtils.toPbWorkerInfo(workerInfo1, false)
- val pbWorkerInfoWithEmptyResource =
PbSerDeUtils.toPbWorkerInfo(workerInfo1, true)
+ val pbWorkerInfo = PbSerDeUtils.toPbWorkerInfo(workerInfo1, false, false)
+ val pbWorkerInfoWithEmptyResource =
PbSerDeUtils.toPbWorkerInfo(workerInfo1, true, false)
val restoredWorkerInfo = PbSerDeUtils.fromPbWorkerInfo(pbWorkerInfo)
val restoredWorkerInfoWithEmptyResource =
PbSerDeUtils.fromPbWorkerInfo(pbWorkerInfoWithEmptyResource)
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 b7eb192e8..6c74606fc 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
@@ -1018,6 +1018,8 @@ private[celeborn] class Master(
// unknown workers will retain in needCheckedWorkerList
needCheckedWorkerList.removeAll(statusSystem.workers)
if (shouldResponse) {
+ // UserResourceConsumption and DiskInfo are eliminated from WorkerInfo
+ // during serialization of HeartbeatFromApplicationResponse
context.reply(HeartbeatFromApplicationResponse(
StatusCode.SUCCESS,
new util.ArrayList(