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 ad13b04f2 [CELEBORN-732] Remove unused RPC ThreadDump &
ThreadDumpResponse
ad13b04f2 is described below
commit ad13b04f2e7efe77ef130af66058f5879ed89029
Author: Angerszhuuuu <[email protected]>
AuthorDate: Wed Jun 28 19:43:39 2023 +0800
[CELEBORN-732] Remove unused RPC ThreadDump & ThreadDumpResponse
### What changes were proposed in this pull request?
Remove unused RPC ThreadDump & ThreadDumpResponse
### Why are the changes needed?
### Does this PR introduce _any_ user-facing change?
### How was this patch tested?
Closes #1645 from AngersZhuuuu/CELEBORN-732.
Authored-by: Angerszhuuuu <[email protected]>
Signed-off-by: Cheng Pan <[email protected]>
---
common/src/main/proto/TransportMessages.proto | 8 ++------
.../common/protocol/message/ControlMessages.scala | 19 -------------------
.../celeborn/service/deploy/worker/Controller.scala | 8 --------
3 files changed, 2 insertions(+), 33 deletions(-)
diff --git a/common/src/main/proto/TransportMessages.proto
b/common/src/main/proto/TransportMessages.proto
index 6ed80b7ff..d1430d72b 100644
--- a/common/src/main/proto/TransportMessages.proto
+++ b/common/src/main/proto/TransportMessages.proto
@@ -55,8 +55,8 @@ enum MessageType {
// SLAVE_LOST_RESPONSE = 34;
GET_WORKER_INFO = 35;
GET_WORKER_INFO_RESPONSE = 36;
- THREAD_DUMP = 37;
- THREAD_DUMP_RESPONSE = 38;
+ // THREAD_DUMP = 37;
+ // THREAD_DUMP_RESPONSE = 38;
REMOVE_EXPIRED_SHUFFLE = 39;
ONE_WAY_MESSAGE_RESPONSE = 40;
CHECK_FOR_WORKER_TIMEOUT = 41;
@@ -390,10 +390,6 @@ message PbGetWorkerInfosResponse {
repeated PbWorkerInfo workerInfos = 2;
}
-message PbThreadDumpResponse {
- string threadDump = 1;
-}
-
message PbCheckForWorkerTimeout {
}
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 3fbff5d81..bafbb2bc0 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
@@ -424,10 +424,6 @@ object ControlMessages extends Logging {
case class GetWorkerInfosResponse(status: StatusCode, workerInfos:
WorkerInfo*) extends Message
- case object ThreadDump extends Message
-
- case class ThreadDumpResponse(threadDump: String) extends Message
-
// TODO change message type to GeneratedMessageV3
def toTransportMessage(message: Any): TransportMessage = message match {
case _: PbCheckForWorkerTimeoutOrBuilder =>
@@ -796,14 +792,6 @@ object ControlMessages extends Logging {
.build().toByteArray
new TransportMessage(MessageType.GET_WORKER_INFO_RESPONSE, payload)
- case ThreadDump =>
- new TransportMessage(MessageType.THREAD_DUMP, null)
-
- case ThreadDumpResponse(threadDump) =>
- val payload = PbThreadDumpResponse.newBuilder()
- .setThreadDump(threadDump).build().toByteArray
- new TransportMessage(MessageType.THREAD_DUMP_RESPONSE, payload)
-
case pb: PbPartitionSplit =>
new TransportMessage(MessageType.PARTITION_SPLIT, pb.toByteArray)
@@ -1100,13 +1088,6 @@ object ControlMessages extends Logging {
pbGetWorkerInfoResponse.getWorkerInfosList.asScala
.map(PbSerDeUtils.fromPbWorkerInfo).toList: _*)
- case THREAD_DUMP =>
- ThreadDump
-
- case THREAD_DUMP_RESPONSE =>
- val pbThreadDumpResponse =
PbThreadDumpResponse.parseFrom(message.getPayload)
- ThreadDumpResponse(pbThreadDumpResponse.getThreadDump)
-
case REMOVE_EXPIRED_SHUFFLE =>
RemoveExpiredShuffle
diff --git
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Controller.scala
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Controller.scala
index 5093c05a1..d1c5984a5 100644
---
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Controller.scala
+++
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Controller.scala
@@ -124,9 +124,6 @@ private[deploy] class Controller(
case GetWorkerInfos =>
handleGetWorkerInfos(context)
- case ThreadDump =>
- handleThreadDump(context)
-
case DestroyWorkerSlots(shuffleKey, masterLocations, slaveLocations) =>
handleDestroy(context, shuffleKey, masterLocations, slaveLocations)
}
@@ -671,9 +668,4 @@ private[deploy] class Controller(
list.add(workerInfo)
context.reply(GetWorkerInfosResponse(StatusCode.SUCCESS,
list.asScala.toList: _*))
}
-
- private def handleThreadDump(context: RpcCallContext): Unit = {
- val threadDump = Utils.getThreadDump()
- context.reply(ThreadDumpResponse(threadDump))
- }
}