This is an automated email from the ASF dual-hosted git repository.
rexxiong 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 e71d912d5 [CELEBORN-1245] Support Celeborn Master(Leader) to manage
workers
e71d912d5 is described below
commit e71d912d5005023f2e843b982c9a40d4b094c129
Author: Shuang <[email protected]>
AuthorDate: Thu Feb 1 09:44:59 2024 +0800
[CELEBORN-1245] Support Celeborn Master(Leader) to manage workers
### What changes were proposed in this pull request?
1. Support Celeborn Master(Leader) to manage workers by sending event when
heartbeat
2. Add Worker Status to Worker then we can know the status of the
workers(such as during decommission...)
3. Add Http interface for master to handleWorkerEvent/getWorkerEvent
### Why are the changes needed?
Currently, we only support managing the status of workers on the worker
side. This pr supports the master to manage the status of all workers. By
sending events such as (Decommission/Graceful/Exit) when heartbeat, workers can
be asynchronously execute the command from master. MeanWhile we can't know what
the worker status during worker decommission so this pr add worker status to
tell the exactly status of the worker.
### Does this PR introduce _any_ user-facing change?
No
### How was this patch tested?
Pass GA
Closes #2255 from RexXiong/CELEBORN-1245.
Authored-by: Shuang <[email protected]>
Signed-off-by: Shuang <[email protected]>
---
.../celeborn/common/util/WorkerStatusUtils.java | 42 ++++++
common/src/main/proto/TransportMessages.proto | 45 ++++++
.../celeborn/common/meta/WorkerEventInfo.java | 76 ++++++++++
.../apache/celeborn/common/meta/WorkerInfo.scala | 10 ++
.../apache/celeborn/common/meta/WorkerStatus.java | 74 ++++++++++
.../common/protocol/message/ControlMessages.scala | 48 ++++++-
.../apache/celeborn/common/util/PbSerDeUtils.scala | 33 ++++-
.../celeborn/common/meta/WorkerInfoSuite.scala | 16 ++-
.../celeborn/common/util/PbSerDeUtilsTest.scala | 19 ++-
docs/monitoring.md | 34 ++---
.../master/clustermeta/AbstractMetaManager.java | 60 +++++++-
.../master/clustermeta/IMetadataHandler.java | 5 +
.../deploy/master/clustermeta/MetaUtil.java | 12 ++
.../clustermeta/SingleMasterMetaManager.java | 8 ++
.../master/clustermeta/ha/HAMasterMetaManager.java | 27 ++++
.../deploy/master/clustermeta/ha/MetaHandler.java | 18 +++
master/src/main/proto/Resource.proto | 33 +++++
.../celeborn/service/deploy/master/Master.scala | 80 ++++++++++-
.../clustermeta/DefaultMetaSystemSuiteJ.java | 7 +
.../ha/RatisMasterStatusSystemSuiteJ.java | 6 +
.../celeborn/server/common/HttpService.scala | 5 +
.../celeborn/server/common/http/HttpEndpoint.scala | 20 +++
.../celeborn/server/common/http/HttpUtils.scala | 2 +
.../server/common/http/HttpUtilsSuite.scala | 9 ++
.../celeborn/service/deploy/worker/Worker.scala | 62 +++++---
.../deploy/worker/WorkerStatusManager.scala | 158 +++++++++++++++++++++
.../deploy/worker/WorkerStatusManagerSuite.scala | 64 +++++++++
27 files changed, 918 insertions(+), 55 deletions(-)
diff --git
a/common/src/main/java/org/apache/celeborn/common/util/WorkerStatusUtils.java
b/common/src/main/java/org/apache/celeborn/common/util/WorkerStatusUtils.java
new file mode 100644
index 000000000..ad495d2f2
--- /dev/null
+++
b/common/src/main/java/org/apache/celeborn/common/util/WorkerStatusUtils.java
@@ -0,0 +1,42 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.celeborn.common.util;
+
+import org.apache.celeborn.common.meta.WorkerEventInfo;
+import org.apache.celeborn.common.meta.WorkerStatus;
+import org.apache.celeborn.common.protocol.PbWorkerStatus;
+
+public class WorkerStatusUtils {
+
+ public static boolean meetFinalState(WorkerEventInfo workerEventInfo,
WorkerStatus workerStatus) {
+ switch (workerEventInfo.getEventType()) {
+ case DecommissionThenIdle:
+ return (workerStatus != null && workerStatus.getState() ==
PbWorkerStatus.State.Idle);
+ case Recommission:
+ case None:
+ return (workerStatus != null && workerStatus.getState() ==
PbWorkerStatus.State.Normal);
+ case Immediately:
+ case Decommission:
+ case Graceful:
+ case UNRECOGNIZED:
+ return false;
+ default:
+ throw new IllegalStateException("Unexpected value: " +
workerEventInfo.getEventType());
+ }
+ }
+}
diff --git a/common/src/main/proto/TransportMessages.proto
b/common/src/main/proto/TransportMessages.proto
index 963bc57c9..534445fb3 100644
--- a/common/src/main/proto/TransportMessages.proto
+++ b/common/src/main/proto/TransportMessages.proto
@@ -96,6 +96,8 @@ enum MessageType {
AUTHENTICATION_INITIATION_RESPONSE = 73;
REGISTER_APPLICATION_REQUEST = 74;
REGISTER_APPLICATION_RESPONSE = 75;
+ WORKER_EVENT_REQUEST = 76;
+ WORKER_EVENT_RESPONSE = 77;
}
enum StreamType {
@@ -103,6 +105,15 @@ enum StreamType {
CreditStream = 1;
}
+enum WorkerEventType {
+ None = 0;
+ Immediately = 1; // Immediately exit
+ Decommission = 2; // From normal to decommission, then exit
+ DecommissionThenIdle = 3; // From normal to decommission and keep as idle
state
+ Graceful = 4; // From normal to graceful, then exit
+ Recommission = 5; // -> From InDecommissionThenIdle/Idle state to Normal
+}
+
message PbStorageInfo {
int32 type = 1;
string mountPoint = 2;
@@ -182,11 +193,44 @@ message PbHeartbeatFromWorker {
map<string, PbResourceConsumption> userResourceConsumption = 9;
map<string, int64> estimatedAppDiskUsage = 10;
bool highWorkload = 11;
+ PbWorkerStatus workerStatus = 12;
+}
+
+message PbWorkerStatus {
+ enum State {
+ Normal = 0;
+ Idle = 1; // Can be recommissioned
+ Exit = 2;
+ InDecommissionThenIdle = 3; // Can be recommissioned
+ InDecommission = 4;
+ InGraceFul = 5;
+ InExit = 6;
+ }
+
+ State state = 1;
+ int64 stateStartTime = 2;
}
message PbHeartbeatFromWorkerResponse {
repeated string expiredShuffleKeys = 1;
bool registered = 2;
+ WorkerEventType workerEventType = 3;
+}
+
+message PbWorkerEventInfo {
+ WorkerEventType workerEventType = 1;
+ int64 eventStartTime = 2;
+}
+
+message PbWorkerEventRequest {
+ repeated PbWorkerInfo workers = 1;
+ WorkerEventType workerEventType = 2;
+ string requestId = 3;
+}
+
+message PbWorkerEventResponse {
+ bool success = 1;
+ string message = 2;
}
message PbRegisterShuffle {
@@ -565,6 +609,7 @@ message PbSnapshotMetaInfo {
map<string, int64> lostWorkers = 12;
repeated PbWorkerInfo shutdownWorkers = 13;
repeated PbWorkerInfo manuallyExcludedWorkers = 14;
+ map<string, PbWorkerEventInfo> workerEventInfos = 15;
}
message PbOpenStream {
diff --git
a/common/src/main/scala/org/apache/celeborn/common/meta/WorkerEventInfo.java
b/common/src/main/scala/org/apache/celeborn/common/meta/WorkerEventInfo.java
new file mode 100644
index 000000000..298c78cbb
--- /dev/null
+++ b/common/src/main/scala/org/apache/celeborn/common/meta/WorkerEventInfo.java
@@ -0,0 +1,76 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.celeborn.common.meta;
+
+import org.apache.celeborn.common.protocol.WorkerEventType;
+
+import java.io.Serializable;
+import java.util.Objects;
+
+public class WorkerEventInfo implements Serializable {
+ private static final long serialVersionUID = 5681914909039445235L;
+ private int eventTypeValue;
+ private long eventStartTime;
+
+ public WorkerEventInfo(int eventTypeValue, long eventStartTime) {
+ this.eventTypeValue = eventTypeValue;
+ this.eventStartTime = eventStartTime;
+ }
+
+ public boolean isSameEvent(int eventTypeValue) {
+ return this.eventTypeValue == eventTypeValue;
+ }
+
+ public WorkerEventType getEventType() {
+ return WorkerEventType.forNumber(eventTypeValue);
+ }
+
+ public long getEventStartTime() {
+ return eventStartTime;
+ }
+
+ public void setEventStartTime(long eventStartTime) {
+ this.eventStartTime = eventStartTime;
+ }
+
+ @Override
+ public boolean equals(Object o) {
+ if (this == o) {
+ return true;
+ }
+ if (!(o instanceof WorkerEventInfo)) {
+ return false;
+ }
+ WorkerEventInfo that = (WorkerEventInfo) o;
+ return eventTypeValue == that.eventTypeValue && getEventStartTime() ==
that.getEventStartTime();
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(eventTypeValue, getEventStartTime());
+ }
+
+ @Override
+ public String toString() {
+ final StringBuilder sb = new StringBuilder("WorkerEventInfo{");
+ sb.append("eventType=").append(getEventType());
+ sb.append(", eventStartTime=").append(eventStartTime);
+ sb.append('}');
+ return sb.toString();
+ }
+}
diff --git
a/common/src/main/scala/org/apache/celeborn/common/meta/WorkerInfo.scala
b/common/src/main/scala/org/apache/celeborn/common/meta/WorkerInfo.scala
index 683a19557..24c5e5a5e 100644
--- a/common/src/main/scala/org/apache/celeborn/common/meta/WorkerInfo.scala
+++ b/common/src/main/scala/org/apache/celeborn/common/meta/WorkerInfo.scala
@@ -41,6 +41,7 @@ class WorkerInfo(
with Logging {
var networkLocation = "/default-rack"
var lastHeartbeat: Long = 0
+ var workerStatus = WorkerStatus.normalWorkerStatus()
val diskInfos =
if (_diskInfos != null) JavaUtils.newConcurrentHashMap[String,
DiskInfo](_diskInfos) else null
val userResourceConsumption =
@@ -142,6 +143,14 @@ class WorkerInfo(
diskInfos.asScala.map(_._2.maxSlots).sum
}
+ def getWorkerStatus(): WorkerStatus = {
+ workerStatus
+ }
+
+ def setWorkerStatus(workerStatus: WorkerStatus): Unit = {
+ this.workerStatus = workerStatus;
+ }
+
def updateDiskMaxSlots(estimatedPartitionSize: Long): Unit =
this.synchronized {
diskInfos.asScala.foreach { case (_, disk) =>
disk.maxSlots_$eq(disk.actualUsableSpace / estimatedPartitionSize)
@@ -230,6 +239,7 @@ class WorkerInfo(
|Disks: $diskInfosString
|UserResourceConsumption: $userResourceConsumptionString
|WorkerRef: $endpoint
+ |WorkerStatus: $workerStatus
|""".stripMargin
}
diff --git
a/common/src/main/scala/org/apache/celeborn/common/meta/WorkerStatus.java
b/common/src/main/scala/org/apache/celeborn/common/meta/WorkerStatus.java
new file mode 100644
index 000000000..a432b540f
--- /dev/null
+++ b/common/src/main/scala/org/apache/celeborn/common/meta/WorkerStatus.java
@@ -0,0 +1,74 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.celeborn.common.meta;
+
+import org.apache.celeborn.common.protocol.PbWorkerStatus;
+
+import java.util.Objects;
+
+public class WorkerStatus {
+ private int stateValue;
+ private long stateStartTime;
+
+ public WorkerStatus(int stateValue, long stateStartTime) {
+ this.stateValue = stateValue;
+ this.stateStartTime = stateStartTime;
+ }
+
+ public int getStateValue() {
+ return stateValue;
+ }
+
+ public PbWorkerStatus.State getState() {
+ return PbWorkerStatus.State.forNumber(stateValue);
+ }
+
+ public long getStateStartTime() {
+ return stateStartTime;
+ }
+
+ public static WorkerStatus normalWorkerStatus() {
+ return new WorkerStatus(PbWorkerStatus.State.Normal.getNumber(),
System.currentTimeMillis());
+ }
+
+ @Override
+ public boolean equals(Object o) {
+ if (this == o) {
+ return true;
+ }
+ if (!(o instanceof WorkerStatus)) {
+ return false;
+ }
+ WorkerStatus that = (WorkerStatus) o;
+ return getStateValue() == that.getStateValue() && getStateStartTime() ==
that.getStateStartTime();
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(getStateValue(), getStateStartTime());
+ }
+
+ @Override
+ public String toString() {
+ final StringBuilder sb = new StringBuilder("WorkerStatus{");
+ sb.append("state=").append(getState());
+ sb.append(", stateStartTime=").append(stateStartTime);
+ sb.append('}');
+ return sb.toString();
+ }
+}
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 aa5b9484a..fed0cd95d 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
@@ -26,7 +26,7 @@ import org.roaringbitmap.RoaringBitmap
import org.apache.celeborn.common.identity.UserIdentifier
import org.apache.celeborn.common.internal.Logging
-import org.apache.celeborn.common.meta.{DiskInfo, WorkerInfo}
+import org.apache.celeborn.common.meta.{DiskInfo, WorkerInfo, WorkerStatus}
import org.apache.celeborn.common.network.protocol.TransportMessage
import org.apache.celeborn.common.protocol._
import org.apache.celeborn.common.protocol.MessageType._
@@ -116,11 +116,14 @@ object ControlMessages extends Logging {
activeShuffleKeys: util.Set[String],
estimatedAppDiskUsage: util.HashMap[String, java.lang.Long],
highWorkload: Boolean,
+ workerStatus: WorkerStatus,
override var requestId: String = ZERO_UUID) extends MasterRequestMessage
case class HeartbeatFromWorkerResponse(
expiredShuffleKeys: util.HashSet[String],
- registered: Boolean) extends MasterMessage
+ registered: Boolean,
+ workerEvent: WorkerEventType = WorkerEventType.None)
+ extends MasterMessage
object RegisterShuffle {
def apply(
@@ -398,6 +401,20 @@ object ControlMessages extends Logging {
.build()
}
+ object WorkerEventRequest {
+ def apply(
+ workers: util.List[WorkerInfo],
+ eventType: String,
+ requestId: String): PbWorkerEventRequest =
+ PbWorkerEventRequest.newBuilder()
+ .setRequestId(requestId)
+ .setWorkerEventType(WorkerEventType.valueOf(eventType))
+ .addAllWorkers(workers.asScala.map { workerInfo =>
+ PbSerDeUtils.toPbWorkerInfo(workerInfo, true)
+ }.toList.asJava)
+ .build()
+ }
+
/**
* ==========================================
* handled by worker
@@ -513,6 +530,7 @@ object ControlMessages extends Logging {
activeShuffleKeys,
estimatedAppDiskUsage,
highWorkload,
+ workerStatus,
requestId) =>
val pbDisks = disks.map(PbSerDeUtils.toPbDiskInfo).asJava
val pbUserResourceConsumption =
@@ -528,14 +546,16 @@ object ControlMessages extends Logging {
.addAllActiveShuffleKeys(activeShuffleKeys)
.putAllEstimatedAppDiskUsage(estimatedAppDiskUsage)
.setHighWorkload(highWorkload)
+ .setWorkerStatus(PbSerDeUtils.toPbWorkerStatus(workerStatus))
.setRequestId(requestId)
.build().toByteArray
new TransportMessage(MessageType.HEARTBEAT_FROM_WORKER, payload)
- case HeartbeatFromWorkerResponse(expiredShuffleKeys, registered) =>
+ case HeartbeatFromWorkerResponse(expiredShuffleKeys, registered,
workerEventType) =>
val payload = PbHeartbeatFromWorkerResponse.newBuilder()
.addAllExpiredShuffleKeys(expiredShuffleKeys)
.setRegistered(registered)
+ .setWorkerEventType(workerEventType)
.build().toByteArray
new TransportMessage(MessageType.HEARTBEAT_FROM_WORKER_RESPONSE, payload)
@@ -747,6 +767,12 @@ object ControlMessages extends Logging {
case pb: PbRemoveWorkersUnavailableInfo =>
new TransportMessage(MessageType.REMOVE_WORKERS_UNAVAILABLE_INFO,
pb.toByteArray)
+ case pb: PbWorkerEventRequest =>
+ new TransportMessage(MessageType.WORKER_EVENT_REQUEST, pb.toByteArray)
+
+ case pb: PbWorkerEventResponse =>
+ new TransportMessage(MessageType.WORKER_EVENT_RESPONSE, pb.toByteArray)
+
case pb: PbRegisterWorkerResponse =>
new TransportMessage(MessageType.REGISTER_WORKER_RESPONSE,
pb.toByteArray)
@@ -910,6 +936,9 @@ object ControlMessages extends Logging {
if (!pbHeartbeatFromWorker.getActiveShuffleKeysList.isEmpty) {
activeShuffleKeys.addAll(pbHeartbeatFromWorker.getActiveShuffleKeysList)
}
+
+ val workerStatus =
PbSerDeUtils.fromPbWorkerStatus(pbHeartbeatFromWorker.getWorkerStatus)
+
HeartbeatFromWorker(
pbHeartbeatFromWorker.getHost,
pbHeartbeatFromWorker.getRpcPort,
@@ -921,6 +950,7 @@ object ControlMessages extends Logging {
activeShuffleKeys,
estimatedAppDiskUsage,
pbHeartbeatFromWorker.getHighWorkload,
+ workerStatus,
pbHeartbeatFromWorker.getRequestId)
case HEARTBEAT_FROM_WORKER_RESPONSE_VALUE =>
@@ -930,7 +960,11 @@ object ControlMessages extends Logging {
if (pbHeartbeatFromWorkerResponse.getExpiredShuffleKeysCount > 0) {
expiredShuffleKeys.addAll(pbHeartbeatFromWorkerResponse.getExpiredShuffleKeysList)
}
- HeartbeatFromWorkerResponse(expiredShuffleKeys,
pbHeartbeatFromWorkerResponse.getRegistered)
+
+ HeartbeatFromWorkerResponse(
+ expiredShuffleKeys,
+ pbHeartbeatFromWorkerResponse.getRegistered,
+ pbHeartbeatFromWorkerResponse.getWorkerEventType)
case REGISTER_SHUFFLE_VALUE =>
PbRegisterShuffle.parseFrom(message.getPayload)
@@ -1083,6 +1117,12 @@ object ControlMessages extends Logging {
case REGISTER_WORKER_RESPONSE_VALUE =>
PbRegisterWorkerResponse.parseFrom(message.getPayload)
+ case WORKER_EVENT_REQUEST_VALUE =>
+ PbWorkerEventRequest.parseFrom(message.getPayload)
+
+ case WORKER_EVENT_RESPONSE_VALUE =>
+ PbWorkerEventResponse.parseFrom(message.getPayload)
+
case RESERVE_SLOTS_VALUE =>
val pbReserveSlots = PbReserveSlots.parseFrom(message.getPayload)
val userIdentifier =
PbSerDeUtils.fromPbUserIdentifier(pbReserveSlots.getUserIdentifier)
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 da50d9965..d7e10fecf 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
@@ -26,7 +26,7 @@ import scala.collection.JavaConverters._
import com.google.protobuf.InvalidProtocolBufferException
import org.apache.celeborn.common.identity.UserIdentifier
-import org.apache.celeborn.common.meta.{AppDiskUsage, AppDiskUsageSnapShot,
DiskFileInfo, DiskInfo, FileInfo, MapFileMeta, ReduceFileMeta, WorkerInfo}
+import org.apache.celeborn.common.meta.{AppDiskUsage, AppDiskUsageSnapShot,
DiskFileInfo, DiskInfo, FileInfo, MapFileMeta, ReduceFileMeta, WorkerEventInfo,
WorkerInfo, WorkerStatus}
import org.apache.celeborn.common.protocol._
import org.apache.celeborn.common.protocol.PartitionLocation.Mode
import
org.apache.celeborn.common.protocol.message.ControlMessages.WorkerResource
@@ -393,7 +393,8 @@ object PbSerDeUtils {
appDiskUsageMetricSnapshots: Array[AppDiskUsageSnapShot],
currentAppDiskUsageMetricsSnapshot: AppDiskUsageSnapShot,
lostWorkers: ConcurrentHashMap[WorkerInfo, java.lang.Long],
- shutdownWorkers: java.util.Set[WorkerInfo]): PbSnapshotMetaInfo = {
+ shutdownWorkers: java.util.Set[WorkerInfo],
+ workerEventInfos: ConcurrentHashMap[WorkerInfo, WorkerEventInfo]):
PbSnapshotMetaInfo = {
val builder = PbSnapshotMetaInfo.newBuilder()
.setEstimatedPartitionSize(estimatedPartitionSize)
.addAllRegisteredShuffle(registeredShuffle)
@@ -414,10 +415,38 @@ object PbSerDeUtils {
case (worker: WorkerInfo, time: java.lang.Long) =>
(worker.toUniqueId(), time)
}.asJava)
.addAllShutdownWorkers(shutdownWorkers.asScala.map(toPbWorkerInfo(_,
true)).asJava)
+ .putAllWorkerEventInfos(workerEventInfos.asScala.map {
+ case (worker, workerEventInfo) =>
+ (worker.toUniqueId(),
PbSerDeUtils.toPbWorkerEventInfo(workerEventInfo))
+ }.asJava)
if (currentAppDiskUsageMetricsSnapshot != null) {
builder.setCurrentAppDiskUsageMetricsSnapshot(
toPbAppDiskUsageSnapshot(currentAppDiskUsageMetricsSnapshot))
}
builder.build()
}
+
+ def toPbWorkerStatus(workerStatus: WorkerStatus): PbWorkerStatus = {
+ PbWorkerStatus.newBuilder()
+ .setState(workerStatus.getState)
+ .setStateStartTime(workerStatus.getStateStartTime)
+ .build()
+ }
+
+ def fromPbWorkerStatus(pbWorkerStatus: PbWorkerStatus): WorkerStatus = {
+ new WorkerStatus(pbWorkerStatus.getState.getNumber,
pbWorkerStatus.getStateStartTime)
+ }
+
+ def toPbWorkerEventInfo(workerEventInfo: WorkerEventInfo): PbWorkerEventInfo
= {
+ PbWorkerEventInfo.newBuilder()
+ .setEventStartTime(workerEventInfo.getEventStartTime)
+ .setWorkerEventType(workerEventInfo.getEventType)
+ .build()
+ }
+
+ def fromPbWorkerEventInfo(pbWorkerEventInfo: PbWorkerEventInfo):
WorkerEventInfo = {
+ new WorkerEventInfo(
+ pbWorkerEventInfo.getWorkerEventType.getNumber,
+ pbWorkerEventInfo.getEventStartTime())
+ }
}
diff --git
a/common/src/test/scala/org/apache/celeborn/common/meta/WorkerInfoSuite.scala
b/common/src/test/scala/org/apache/celeborn/common/meta/WorkerInfoSuite.scala
index a435f35e3..1691c91ba 100644
---
a/common/src/test/scala/org/apache/celeborn/common/meta/WorkerInfoSuite.scala
+++
b/common/src/test/scala/org/apache/celeborn/common/meta/WorkerInfoSuite.scala
@@ -292,10 +292,18 @@ class WorkerInfoSuite extends CelebornFunSuite {
|WorkerRef: null
|""".stripMargin;
- assertEquals(exp1,
worker1.toString.replaceAll("HeartbeatElapsedSeconds:.*\n", ""))
- assertEquals(exp2,
worker2.toString.replaceAll("HeartbeatElapsedSeconds:.*\n", ""))
- assertEquals(exp3,
worker3.toString.replaceAll("HeartbeatElapsedSeconds:.*\n", ""))
- assertEquals(exp4,
worker4.toString.replaceAll("HeartbeatElapsedSeconds:.*\n", ""))
+ assertEquals(
+ exp1,
+
worker1.toString.replaceAll("(HeartbeatElapsedSeconds|WorkerStatus):.*\n", ""))
+ assertEquals(
+ exp2,
+
worker2.toString.replaceAll("(HeartbeatElapsedSeconds|WorkerStatus):.*\n", ""))
+ assertEquals(
+ exp3,
+
worker3.toString.replaceAll("(HeartbeatElapsedSeconds|WorkerStatus):.*\n", ""))
+ assertEquals(
+ exp4,
+
worker4.toString.replaceAll("(HeartbeatElapsedSeconds|WorkerStatus):.*\n", ""))
} finally {
if (null != rpcEnv) {
rpcEnv.shutdown()
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 e374cab93..de475b8ee 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
@@ -22,8 +22,9 @@ import java.util
import org.apache.celeborn.CelebornFunSuite
import org.apache.celeborn.common.identity.UserIdentifier
-import org.apache.celeborn.common.meta.{DeviceInfo, DiskFileInfo, DiskInfo,
FileInfo, ReduceFileMeta, WorkerInfo}
+import org.apache.celeborn.common.meta.{DeviceInfo, DiskFileInfo, DiskInfo,
FileInfo, ReduceFileMeta, WorkerEventInfo, WorkerInfo, WorkerStatus}
import org.apache.celeborn.common.protocol.{PartitionLocation, StorageInfo}
+import org.apache.celeborn.common.protocol.PartitionLocation
import
org.apache.celeborn.common.protocol.message.ControlMessages.WorkerResource
import org.apache.celeborn.common.quota.ResourceConsumption
@@ -113,6 +114,9 @@ class PbSerDeUtilsTest extends CelebornFunSuite {
workerInfo1,
(util.Arrays.asList(partitionLocation1),
util.Arrays.asList(partitionLocation2)))
+ val workerEventInfo = new WorkerEventInfo(1, System.currentTimeMillis());
+ val workerStatus = new WorkerStatus(1, System.currentTimeMillis());
+
test("fromAndToPbSortedShuffleFileSet") {
val pbFileSet = PbSerDeUtils.toPbSortedShuffleFileSet(fileSet)
val restoredFileSet = PbSerDeUtils.fromPbSortedShuffleFileSet(pbFileSet)
@@ -237,4 +241,17 @@ class PbSerDeUtilsTest extends CelebornFunSuite {
assert(restoredPartitionLocation4.getStorageInfo.equals(partitionLocation4.getStorageInfo))
}
+ test("fromAndToPbWorkerEventInfo") {
+ val pbWorkerEventInfo = PbSerDeUtils.toPbWorkerEventInfo(workerEventInfo)
+ val restoredWorkerEventInfo =
PbSerDeUtils.fromPbWorkerEventInfo(pbWorkerEventInfo)
+
+ assert(restoredWorkerEventInfo.equals(workerEventInfo))
+ }
+
+ test("fromAndToPbWorkerStatus") {
+ val pbWorkerStatus = PbSerDeUtils.toPbWorkerStatus(workerStatus)
+ val restoredWorkerStatus = PbSerDeUtils.fromPbWorkerStatus(pbWorkerStatus)
+
+ assert(restoredWorkerStatus.equals(workerStatus))
+ }
}
diff --git a/docs/monitoring.md b/docs/monitoring.md
index f94997f94..0f206a558 100644
--- a/docs/monitoring.md
+++ b/docs/monitoring.md
@@ -344,22 +344,24 @@ API path listed as below:
#### Master
-| Path | Meaning
|
-|------------------------------------------------------|----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|
-| /metrics/prometheus | List the metrics data
in prometheus format of the master.(The url path is defined by configure
`celeborn.metrics.prometheus.path`.)
|
-| /conf | List the conf setting
of the master.
|
-| /masterGroupInfo | List master group
information of the service. It will list all master's LEADER, FOLLOWER
information.
|
-| /workerInfo | List worker
information of the service. It will list all registered workers 's information.
|
-| /lostWorkers | List all lost workers
of the master.
|
-| /excludedWorkers | List all excluded
workers of the master.
|
-| /shutdownWorkers | List all shutdown
workers of the master.
|
-| /threadDump | List the current
thread dump of the master.
|
-| /hostnames | List all running
application's LifecycleManager's hostnames of the cluster.
|
-| /applications | List all running
application's ids of the cluster.
|
-| /shuffles | List all running
shuffle keys of the service. It will return all running shuffle's key of the
cluster.
|
-| /listTopDiskUsedApps | List the top disk
usage application ids. It will return the top disk usage application ids for
the cluster.
|
-| /exclude?add=${ADD_WORKERS}&remove=${REMOVE_WORKERS} | Excluded workers of
the master add or remove the worker manually given worker id. The parameter add
or remove specifies the excluded workers to add or remove, which value is
separated by commas. |
-| /help | List the available
API providers of the master.
|
+| Path | Meaning
|
+|-------------------------------------------------------------|-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|
+| /metrics/prometheus | List the
metrics data in prometheus format of the master.(The url path is defined by
configure `celeborn.metrics.prometheus.path`.)
|
+| /conf | List the conf
setting of the master.
|
+| /masterGroupInfo | List master
group information of the service. It will list all master's LEADER, FOLLOWER
information.
|
+| /workerInfo | List worker
information of the service. It will list all registered workers 's information.
|
+| /lostWorkers | List all lost
workers of the master.
|
+| /excludedWorkers | List all
excluded workers of the master.
|
+| /shutdownWorkers | List all
shutdown workers of the master.
|
+| /threadDump | List the
current thread dump of the master.
|
+| /hostnames | List all
running application's LifecycleManager's hostnames of the cluster.
|
+| /applications | List all
running application's ids of the cluster.
|
+| /shuffles | List all
running shuffle keys of the service. It will return all running shuffle's key
of the cluster.
|
+| /listTopDiskUsedApps | List the top
disk usage application ids. It will return the top disk usage application ids
for the cluster.
|
+| /exclude?add=${ADD_WORKERS}&remove=${REMOVE_WORKERS} | Excluded
workers of the master add or remove the worker manually given worker id. The
parameter add or remove specifies the excluded workers to add or remove, which
value is separated by commas. |
+| /sendWorkerEvent?type=${WorkerEventType}&workers=${WORKERS} | For
Master(Leader) can send worker event to manager workers. Legal
`WorkerEventType` are 'None', 'Immediately', 'Decommission',
'DecommissionThenIdle', 'Graceful', 'Recommission', and the parameter workers
is separated by commas. |
+| /workerEventInfo | List all
worker event infos of the master.
|
+| /help | List the
available API providers of the master.
|
#### Worker
diff --git
a/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/AbstractMetaManager.java
b/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/AbstractMetaManager.java
index 8f468d99e..abe19ec13 100644
---
a/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/AbstractMetaManager.java
+++
b/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/AbstractMetaManager.java
@@ -45,13 +45,17 @@ import org.apache.celeborn.common.meta.AppDiskUsageMetric;
import org.apache.celeborn.common.meta.AppDiskUsageSnapShot;
import org.apache.celeborn.common.meta.DiskInfo;
import org.apache.celeborn.common.meta.DiskStatus;
+import org.apache.celeborn.common.meta.WorkerEventInfo;
import org.apache.celeborn.common.meta.WorkerInfo;
+import org.apache.celeborn.common.meta.WorkerStatus;
import org.apache.celeborn.common.protocol.PbSnapshotMetaInfo;
+import org.apache.celeborn.common.protocol.PbWorkerStatus;
import org.apache.celeborn.common.quota.ResourceConsumption;
import org.apache.celeborn.common.rpc.RpcEnv;
import org.apache.celeborn.common.util.JavaUtils;
import org.apache.celeborn.common.util.PbSerDeUtils;
import org.apache.celeborn.common.util.Utils;
+import org.apache.celeborn.common.util.WorkerStatusUtils;
import org.apache.celeborn.service.deploy.master.network.CelebornRackResolver;
public abstract class AbstractMetaManager implements IMetadataHandler {
@@ -62,6 +66,8 @@ public abstract class AbstractMetaManager implements
IMetadataHandler {
public final Set<String> hostnameSet = ConcurrentHashMap.newKeySet();
public final ArrayList<WorkerInfo> workers = new ArrayList<>();
public final ConcurrentHashMap<WorkerInfo, Long> lostWorkers =
JavaUtils.newConcurrentHashMap();
+ public final ConcurrentHashMap<WorkerInfo, WorkerEventInfo> workerEventInfos
=
+ JavaUtils.newConcurrentHashMap();
public final ConcurrentHashMap<String, Long> appHeartbeatTime =
JavaUtils.newConcurrentHashMap();
public final Set<WorkerInfo> excludedWorkers = ConcurrentHashMap.newKeySet();
public final Set<WorkerInfo> manuallyExcludedWorkers =
ConcurrentHashMap.newKeySet();
@@ -149,6 +155,7 @@ public abstract class AbstractMetaManager implements
IMetadataHandler {
if (lostWorkers.containsKey(workerInfo)) {
lostWorkers.remove(workerInfo);
shutdownWorkers.remove(workerInfo);
+ workerEventInfos.remove(workerInfo);
}
}
}
@@ -164,6 +171,7 @@ public abstract class AbstractMetaManager implements
IMetadataHandler {
Map<UserIdentifier, ResourceConsumption> userResourceConsumption,
Map<String, Long> estimatedAppDiskUsage,
long time,
+ WorkerStatus workerStatus,
boolean highWorkload) {
WorkerInfo worker =
new WorkerInfo(
@@ -178,8 +186,19 @@ public abstract class AbstractMetaManager implements
IMetadataHandler {
info.updateThenGetUserResourceConsumption(userResourceConsumption);
availableSlots.set(info.totalAvailableSlots());
info.lastHeartbeat_$eq(time);
+ info.setWorkerStatus(workerStatus);
});
}
+
+ WorkerEventInfo workerEventInfo = workerEventInfos.get(worker);
+ if (workerEventInfo != null
+ && WorkerStatusUtils.meetFinalState(workerEventInfo, workerStatus)) {
+ workerEventInfos.remove(worker);
+ if (workerStatus.getState() == PbWorkerStatus.State.Normal) {
+ shutdownWorkers.remove(worker);
+ }
+ }
+
appDiskUsageMetric.update(estimatedAppDiskUsage);
// If using HDFSONLY mode, workers with empty disks should not be put into
excluded worker list.
long healthyDiskNum =
@@ -215,6 +234,7 @@ public abstract class AbstractMetaManager implements
IMetadataHandler {
shutdownWorkers.remove(workerInfo);
lostWorkers.remove(workerInfo);
excludedWorkers.remove(workerInfo);
+ workerEventInfos.remove(workerInfo);
}
}
@@ -240,7 +260,8 @@ public abstract class AbstractMetaManager implements
IMetadataHandler {
appDiskUsageMetric.snapShots(),
appDiskUsageMetric.currentSnapShot().get(),
lostWorkers,
- shutdownWorkers)
+ shutdownWorkers,
+ workerEventInfos)
.toByteArray();
Files.write(file.toPath(), snapshotBytes);
}
@@ -304,6 +325,15 @@ public abstract class AbstractMetaManager implements
IMetadataHandler {
.getLostWorkersMap()
.forEach((key, value) ->
lostWorkers.put(WorkerInfo.fromUniqueId(key), value));
+ snapshotMetaInfo
+ .getWorkerEventInfosMap()
+ .entrySet()
+ .forEach(
+ entry ->
+ workerEventInfos.put(
+ WorkerInfo.fromUniqueId(entry.getKey()),
+ PbSerDeUtils.fromPbWorkerEventInfo(entry.getValue())));
+
shutdownWorkers.addAll(
snapshotMetaInfo.getShutdownWorkersList().stream()
.map(PbSerDeUtils::fromPbWorkerInfo)
@@ -345,6 +375,7 @@ public abstract class AbstractMetaManager implements
IMetadataHandler {
workerLostEvents.clear();
partitionTotalWritten.reset();
partitionTotalFileCount.reset();
+ workerEventInfos.clear();
}
public void updateMetaByReportWorkerUnavailable(List<WorkerInfo>
failedWorkers) {
@@ -353,6 +384,25 @@ public abstract class AbstractMetaManager implements
IMetadataHandler {
}
}
+ public void updateWorkerEventMeta(int workerEventTypeValue, List<WorkerInfo>
workerInfoList) {
+ long eventTime = System.currentTimeMillis();
+ ResourceProtos.WorkerEventType eventType =
+ ResourceProtos.WorkerEventType.forNumber(workerEventTypeValue);
+ synchronized (this.workers) {
+ for (WorkerInfo workerInfo : workerInfoList) {
+ WorkerEventInfo workerEventInfo = workerEventInfos.get(workerInfo);
+ LOG.info("Received worker event: {} for worker: {}", eventType,
workerInfo.toUniqueId());
+ if (workerEventInfo == null ||
!workerEventInfo.isSameEvent(eventType.getNumber())) {
+ if (eventType == ResourceProtos.WorkerEventType.None) {
+ workerEventInfos.remove(workerInfo);
+ } else {
+ workerEventInfos.put(workerInfo, new
WorkerEventInfo(eventType.getNumber(), eventTime));
+ }
+ }
+ }
+ }
+ }
+
public void updatePartitionSize() {
long oldEstimatedPartitionSize = estimatedPartitionSize;
long tmpTotalWritten = partitionTotalWritten.sumThenReset();
@@ -376,4 +426,12 @@ public abstract class AbstractMetaManager implements
IMetadataHandler {
!excludedWorkers.contains(worker) &&
!manuallyExcludedWorkers.contains(worker))
.forEach(workerInfo ->
workerInfo.updateDiskMaxSlots(estimatedPartitionSize));
}
+
+ public boolean isWorkerAvailable(WorkerInfo workerInfo) {
+ return !excludedWorkers.contains(workerInfo)
+ && !shutdownWorkers.contains(workerInfo)
+ && !manuallyExcludedWorkers.contains(workerInfo)
+ && (!workerEventInfos.containsKey(workerInfo)
+ && workerInfo.getWorkerStatus().getState() ==
PbWorkerStatus.State.Normal);
+ }
}
diff --git
a/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/IMetadataHandler.java
b/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/IMetadataHandler.java
index a51135738..04fd294f4 100644
---
a/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/IMetadataHandler.java
+++
b/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/IMetadataHandler.java
@@ -23,6 +23,7 @@ import java.util.Map;
import org.apache.celeborn.common.identity.UserIdentifier;
import org.apache.celeborn.common.meta.DiskInfo;
import org.apache.celeborn.common.meta.WorkerInfo;
+import org.apache.celeborn.common.meta.WorkerStatus;
import org.apache.celeborn.common.quota.ResourceConsumption;
public interface IMetadataHandler {
@@ -61,6 +62,7 @@ public interface IMetadataHandler {
Map<String, Long> estimatedAppDiskUsage,
long time,
boolean highWorkload,
+ WorkerStatus workerStatus,
String requestId);
void handleRegisterWorker(
@@ -75,5 +77,8 @@ public interface IMetadataHandler {
void handleReportWorkerUnavailable(List<WorkerInfo> failedNodes, String
requestId);
+ void handleWorkerEvent(
+ int workerEventTypeValue, List<WorkerInfo> workerInfoList, String
requestId);
+
void handleUpdatePartitionSize();
}
diff --git
a/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/MetaUtil.java
b/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/MetaUtil.java
index 03e1d21e3..6ec1f5068 100644
---
a/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/MetaUtil.java
+++
b/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/MetaUtil.java
@@ -24,6 +24,7 @@ import org.apache.celeborn.common.identity.UserIdentifier;
import org.apache.celeborn.common.identity.UserIdentifier$;
import org.apache.celeborn.common.meta.DiskInfo;
import org.apache.celeborn.common.meta.WorkerInfo;
+import org.apache.celeborn.common.meta.WorkerStatus;
import org.apache.celeborn.common.protocol.StorageInfo;
import org.apache.celeborn.common.quota.ResourceConsumption;
import org.apache.celeborn.common.util.Utils;
@@ -120,4 +121,15 @@ public class MetaUtil {
.build()));
return map;
}
+
+ public static ResourceProtos.WorkerStatus toPbWorkerStatus(WorkerStatus
workerStatus) {
+ return ResourceProtos.WorkerStatus.newBuilder()
+
.setState(ResourceProtos.WorkerStatus.State.forNumber(workerStatus.getStateValue()))
+ .setStateStartTime(workerStatus.getStateStartTime())
+ .build();
+ }
+
+ public static WorkerStatus fromPbWorkerStatus(ResourceProtos.WorkerStatus
workerStatus) {
+ return new WorkerStatus(workerStatus.getState().getNumber(),
workerStatus.getStateStartTime());
+ }
}
diff --git
a/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/SingleMasterMetaManager.java
b/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/SingleMasterMetaManager.java
index 1aea3841e..e8a993ef0 100644
---
a/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/SingleMasterMetaManager.java
+++
b/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/SingleMasterMetaManager.java
@@ -28,6 +28,7 @@ import org.apache.celeborn.common.identity.UserIdentifier;
import org.apache.celeborn.common.meta.AppDiskUsageMetric;
import org.apache.celeborn.common.meta.DiskInfo;
import org.apache.celeborn.common.meta.WorkerInfo;
+import org.apache.celeborn.common.meta.WorkerStatus;
import org.apache.celeborn.common.quota.ResourceConsumption;
import org.apache.celeborn.common.rpc.RpcEnv;
import org.apache.celeborn.service.deploy.master.network.CelebornRackResolver;
@@ -110,6 +111,7 @@ public class SingleMasterMetaManager extends
AbstractMetaManager {
Map<String, Long> estimatedAppDiskUsage,
long time,
boolean highWorkload,
+ WorkerStatus workerStatus,
String requestId) {
updateWorkerHeartbeatMeta(
host,
@@ -121,6 +123,7 @@ public class SingleMasterMetaManager extends
AbstractMetaManager {
userResourceConsumption,
estimatedAppDiskUsage,
time,
+ workerStatus,
highWorkload);
}
@@ -144,6 +147,11 @@ public class SingleMasterMetaManager extends
AbstractMetaManager {
}
@Override
+ public void handleWorkerEvent(
+ int workerEventTypeValue, List<WorkerInfo> workerInfoList, String
requestId) {
+ updateWorkerEventMeta(workerEventTypeValue, workerInfoList);
+ }
+
public void handleUpdatePartitionSize() {
updatePartitionSize();
}
diff --git
a/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/HAMasterMetaManager.java
b/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/HAMasterMetaManager.java
index 16fd84b9b..d772d6b1b 100644
---
a/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/HAMasterMetaManager.java
+++
b/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/HAMasterMetaManager.java
@@ -31,6 +31,7 @@ import org.apache.celeborn.common.identity.UserIdentifier;
import org.apache.celeborn.common.meta.AppDiskUsageMetric;
import org.apache.celeborn.common.meta.DiskInfo;
import org.apache.celeborn.common.meta.WorkerInfo;
+import org.apache.celeborn.common.meta.WorkerStatus;
import org.apache.celeborn.common.quota.ResourceConsumption;
import org.apache.celeborn.common.rpc.RpcEnv;
import
org.apache.celeborn.service.deploy.master.clustermeta.AbstractMetaManager;
@@ -253,6 +254,7 @@ public class HAMasterMetaManager extends
AbstractMetaManager {
Map<String, Long> estimatedAppDiskUsage,
long time,
boolean highWorkload,
+ WorkerStatus workerStatus,
String requestId) {
try {
ratisServer.submitRequest(
@@ -270,6 +272,7 @@ public class HAMasterMetaManager extends
AbstractMetaManager {
.putAllUserResourceConsumption(
MetaUtil.toPbUserResourceConsumption(userResourceConsumption))
.putAllEstimatedAppDiskUsage(estimatedAppDiskUsage)
+ .setWorkerStatus(MetaUtil.toPbWorkerStatus(workerStatus))
.setTime(time)
.setHighWorkload(highWorkload)
.build())
@@ -333,6 +336,30 @@ public class HAMasterMetaManager extends
AbstractMetaManager {
}
}
+ @Override
+ public void handleWorkerEvent(
+ int workerEventTypeValue, List<WorkerInfo> workerInfoList, String
requestId) {
+ try {
+ List<ResourceProtos.WorkerAddress> addrs =
+
workerInfoList.stream().map(MetaUtil::infoToAddr).collect(Collectors.toList());
+ ratisServer.submitRequest(
+ ResourceRequest.newBuilder()
+ .setCmdType(Type.WorkerEvent)
+ .setRequestId(requestId)
+ .setWorkerEventRequest(
+ ResourceProtos.WorkerEventRequest.newBuilder()
+ .setWorkerEventType(
+
ResourceProtos.WorkerEventType.forNumber(workerEventTypeValue))
+ .addAllWorkerAddress(addrs)
+ .build())
+ .build());
+ } catch (CelebornRuntimeException e) {
+ LOG.error(
+ "Handle worker event {} failure for {} failed!",
workerEventTypeValue, workerInfoList, e);
+ throw e;
+ }
+ }
+
@Override
public void handleUpdatePartitionSize() {
try {
diff --git
a/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/MetaHandler.java
b/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/MetaHandler.java
index ae660ed17..d9905dcd0 100644
---
a/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/MetaHandler.java
+++
b/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/MetaHandler.java
@@ -29,6 +29,7 @@ import org.apache.celeborn.common.CelebornConf;
import org.apache.celeborn.common.identity.UserIdentifier;
import org.apache.celeborn.common.meta.DiskInfo;
import org.apache.celeborn.common.meta.WorkerInfo;
+import org.apache.celeborn.common.meta.WorkerStatus;
import org.apache.celeborn.common.quota.ResourceConsumption;
import org.apache.celeborn.service.deploy.master.clustermeta.MetaUtil;
import org.apache.celeborn.service.deploy.master.clustermeta.ResourceProtos;
@@ -98,6 +99,7 @@ public class MetaHandler {
Map<String, DiskInfo> diskInfos;
Map<UserIdentifier, ResourceConsumption> userResourceConsumption;
Map<String, Long> estimatedAppDiskUsage = new HashMap<>();
+ WorkerStatus workerStatus;
switch (cmdType) {
case RequestSlots:
shuffleKey = request.getRequestSlotsRequest().getShuffleKey();
@@ -175,6 +177,13 @@ public class MetaHandler {
request.getWorkerHeartbeatRequest().getEstimatedAppDiskUsageMap());
replicatePort =
request.getWorkerHeartbeatRequest().getReplicatePort();
boolean highWorkload =
request.getWorkerHeartbeatRequest().getHighWorkload();
+ if (request.getWorkerHeartbeatRequest().hasWorkerStatus()) {
+ workerStatus =
+
MetaUtil.fromPbWorkerStatus(request.getWorkerHeartbeatRequest().getWorkerStatus());
+ } else {
+ workerStatus = WorkerStatus.normalWorkerStatus();
+ }
+
LOG.debug(
"Handle worker heartbeat for {} {} {} {} {} {} {}",
host,
@@ -194,6 +203,7 @@ public class MetaHandler {
userResourceConsumption,
estimatedAppDiskUsage,
request.getWorkerHeartbeatRequest().getTime(),
+ workerStatus,
highWorkload);
break;
@@ -246,6 +256,14 @@ public class MetaHandler {
metaSystem.removeWorkersUnavailableInfoMeta(unavailableWorkers);
break;
+ case WorkerEvent:
+ List<ResourceProtos.WorkerAddress> workerAddresses =
+
request.getRemoveWorkersUnavailableInfoRequest().getUnavailableList();
+ List<WorkerInfo> workerInfoList =
+
workerAddresses.stream().map(MetaUtil::addrToInfo).collect(Collectors.toList());
+ metaSystem.updateWorkerEventMeta(
+
request.getWorkerEventRequest().getWorkerEventType().getNumber(),
workerInfoList);
+
default:
throw new IOException("Can not parse this command!" + request);
}
diff --git a/master/src/main/proto/Resource.proto
b/master/src/main/proto/Resource.proto
index 78dc477bd..9ee195283 100644
--- a/master/src/main/proto/Resource.proto
+++ b/master/src/main/proto/Resource.proto
@@ -37,6 +37,17 @@ enum Type {
WorkerRemove = 22;
RemoveWorkersUnavailableInfo = 23;
WorkerExclude = 24;
+ WorkerEvent = 25;
+}
+
+enum WorkerEventType {
+ // refer to TransportMessages.proto.PbWorkerEvent.EventType
+ None = 0;
+ Immediately = 1; // Immediately exit
+ Decommission = 2; // from normal to decommission, then exit
+ DecommissionThenIdle = 3; // from normal to decommission and keep alive
+ Graceful = 4; // from normal to graceful, then exit
+ Recommission = 5; // -> from InDecommissionThenIdle/Idle state to Normal
}
message ResourceRequest {
@@ -57,6 +68,7 @@ message ResourceRequest {
optional WorkerRemoveRequest workerRemoveRequest = 19;
optional RemoveWorkersUnavailableInfoRequest
removeWorkersUnavailableInfoRequest = 20;
optional WorkerExcludeRequest workerExcludeRequest = 21;
+ optional WorkerEventRequest workerEventRequest = 22;
}
message DiskInfo {
@@ -132,6 +144,22 @@ message WorkerHeartbeatRequest {
map<string, ResourceConsumption> userResourceConsumption = 8;
map<string, int64> estimatedAppDiskUsage = 9;
required bool highWorkload = 10;
+ optional WorkerStatus workerStatus = 11;
+}
+
+message WorkerStatus {
+ enum State {
+ Normal = 0;
+ Idle = 1;
+ Exit = 2;
+ InDecommissionThenIdle = 3; // -> Idle
+ InDecommission = 4; // -> Exit
+ InGraceFul = 5; // -> Exit
+ InExit = 6;
+ }
+
+ required State state = 1;
+ required int64 stateStartTime = 2;
}
message RegisterWorkerRequest {
@@ -152,6 +180,11 @@ message RemoveWorkersUnavailableInfoRequest {
repeated WorkerAddress unavailable = 1;
}
+message WorkerEventRequest {
+ repeated WorkerAddress workerAddress = 1;
+ required WorkerEventType workerEventType = 2;
+}
+
message WorkerAddress {
required string host = 1;
required int32 rpcPort = 2;
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 bcc2b0eb8..4bb789822 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
@@ -35,7 +35,7 @@ import org.apache.celeborn.common.CelebornConf
import org.apache.celeborn.common.client.MasterClient
import org.apache.celeborn.common.identity.UserIdentifier
import org.apache.celeborn.common.internal.Logging
-import org.apache.celeborn.common.meta.{DiskInfo, WorkerInfo}
+import org.apache.celeborn.common.meta.{DiskInfo, WorkerInfo, WorkerStatus}
import org.apache.celeborn.common.metrics.MetricsSystem
import org.apache.celeborn.common.metrics.source.{JVMCPUSource, JVMSource,
ResourceConsumptionSource, SystemMiscSource, ThreadPoolSource}
import org.apache.celeborn.common.protocol._
@@ -423,6 +423,7 @@ private[celeborn] class Master(
activeShuffleKey,
estimatedAppDiskUsage,
highWorkload,
+ workerStatus,
requestId) =>
logDebug(s"Received heartbeat from" +
s" worker $host:$rpcPort:$pushPort:$fetchPort:$replicatePort with
$disks.")
@@ -440,6 +441,7 @@ private[celeborn] class Master(
activeShuffleKey,
estimatedAppDiskUsage,
highWorkload,
+ workerStatus,
requestId))
case ReportWorkerUnavailable(failedWorkers: util.List[WorkerInfo],
requestId: String) =>
@@ -477,6 +479,17 @@ private[celeborn] class Master(
case _: PbCheckWorkersAvailable =>
executeWithLeaderChecker(context, handleCheckWorkersAvailable(context))
+
+ case pb: PbWorkerEventRequest =>
+ val workers = new util.ArrayList[WorkerInfo](pb.getWorkersList
+ .asScala.map(PbSerDeUtils.fromPbWorkerInfo).toList.asJava)
+ executeWithLeaderChecker(
+ context,
+ handleWorkerEvent(
+ pb.getRequestId,
+ pb.getWorkerEventType.getNumber,
+ workers,
+ context))
}
private def timeoutDeadWorkers(): Unit = {
@@ -556,6 +569,7 @@ private[celeborn] class Master(
activeShuffleKeys: util.Set[String],
estimatedAppDiskUsage: util.HashMap[String, java.lang.Long],
highWorkload: Boolean,
+ workerStatus: WorkerStatus,
requestId: String): Unit = {
val targetWorker = new WorkerInfo(host, rpcPort, pushPort, fetchPort,
replicatePort)
val registered = workersSnapShot.asScala.contains(targetWorker)
@@ -574,6 +588,7 @@ private[celeborn] class Master(
estimatedAppDiskUsage,
System.currentTimeMillis(),
highWorkload,
+ workerStatus,
requestId)
}
@@ -585,7 +600,18 @@ private[celeborn] class Master(
expiredShuffleKeys.add(shuffleKey)
}
}
- context.reply(HeartbeatFromWorkerResponse(expiredShuffleKeys, registered))
+
+ val workerEventInfo = statusSystem.workerEventInfos.get(targetWorker)
+ if (workerEventInfo == null) {
+ context.reply(HeartbeatFromWorkerResponse(
+ expiredShuffleKeys,
+ registered))
+ } else {
+ context.reply(HeartbeatFromWorkerResponse(
+ expiredShuffleKeys,
+ registered,
+ workerEventInfo.getEventType))
+ }
}
private def handleWorkerExclude(
@@ -947,11 +973,19 @@ private[celeborn] class Master(
context.reply(CheckWorkersAvailableResponse(!workersAvailable().isEmpty))
}
+ private def handleWorkerEvent(
+ requestId: String,
+ workerEventTypeValue: Int,
+ workers: util.List[WorkerInfo],
+ context: RpcCallContext): Unit = {
+ statusSystem.handleWorkerEvent(workerEventTypeValue, workers, requestId)
+ context.reply(PbWorkerEventResponse.newBuilder().setSuccess(true).build())
+ }
+
private def workersAvailable(
tmpExcludedWorkerList: Set[WorkerInfo] = Set.empty):
util.List[WorkerInfo] = {
workersSnapShot.asScala.filter { w =>
- !statusSystem.excludedWorkers.contains(w) &&
!statusSystem.manuallyExcludedWorkers.contains(
- w) && !statusSystem.shutdownWorkers.contains(w) &&
!tmpExcludedWorkerList.contains(w)
+ statusSystem.isWorkerAvailable(w) && !tmpExcludedWorkerList.contains(w)
}.asJava
}
@@ -966,6 +1000,35 @@ private[celeborn] class Master(
workersSnapShot.asScala.mkString("\n")
}
+ override def handleWorkerEvent(workerEventType: String, workers: String):
String = {
+ val sb = new StringBuilder
+ if (workerEventType.isEmpty || workers.isEmpty) {
+ return sb.append(
+ s"handle eventType failed as eventType: $workerEventType or workers:
$workers has empty value").toString()
+ }
+
+ sb.append("============================ Handle Worker Event
=============================\n")
+ val workerArray = workers.split(",").filter(_.nonEmpty)
+ try {
+ val workerEventResponse =
self.askSync[PbWorkerEventResponse](WorkerEventRequest(
+ workerArray.map(WorkerInfo.fromUniqueId).toList.asJava,
+ workerEventType,
+ MasterClient.genRequestId()))
+ if (workerEventResponse.getSuccess) {
+ sb.append(s"handle $workerEventType for ${workerArray.mkString(",")}
successfully")
+ } else {
+ sb.append(s"handle $workerEventType for ${workerArray.mkString(",")}
failed")
+ }
+ } catch {
+ case e: Throwable =>
+ val message =
+ s"handle $workerEventType for ${workerArray.mkString(",")} failed,
message: ${e.getMessage}"
+ logError(message, e)
+ sb.append(message)
+ }
+ sb.append("\n").toString()
+ }
+
override def getWorkerInfo: String = {
val sb = new StringBuilder
sb.append("====================== Workers Info in Master
=========================\n")
@@ -1115,6 +1178,15 @@ private[celeborn] class Master(
}
}
+ override def getWorkerEventInfo(): String = {
+ val sb = new StringBuilder
+ sb.append("======================= Workers Event in Master
========================\n")
+ statusSystem.workerEventInfos.asScala.foreach { case (worker,
workerEventInfo) =>
+ sb.append(s"${worker.toUniqueId().padTo(50, "
").mkString}$workerEventInfo\n")
+ }
+ sb.toString()
+ }
+
override def initialize(): Unit = {
super.initialize()
logInfo("Master started.")
diff --git
a/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/DefaultMetaSystemSuiteJ.java
b/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/DefaultMetaSystemSuiteJ.java
index 83f5f7232..b1ec73d3e 100644
---
a/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/DefaultMetaSystemSuiteJ.java
+++
b/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/DefaultMetaSystemSuiteJ.java
@@ -35,6 +35,7 @@ import org.apache.celeborn.common.client.MasterClient;
import org.apache.celeborn.common.identity.UserIdentifier;
import org.apache.celeborn.common.meta.DiskInfo;
import org.apache.celeborn.common.meta.WorkerInfo;
+import org.apache.celeborn.common.meta.WorkerStatus;
import org.apache.celeborn.common.quota.ResourceConsumption;
import org.apache.celeborn.common.rpc.RpcEndpointAddress;
import org.apache.celeborn.common.rpc.RpcEndpointRef;
@@ -78,6 +79,8 @@ public class DefaultMetaSystemSuiteJ {
private static final Map<UserIdentifier, ResourceConsumption>
userResourceConsumption3 =
new HashMap<>();
+ private static final WorkerStatus workerStatus =
WorkerStatus.normalWorkerStatus();
+
@Before
public void setUp() {
when(mockRpcEnv.setupEndpointRef(any(), any())).thenReturn(dummyRef);
@@ -544,6 +547,7 @@ public class DefaultMetaSystemSuiteJ {
new HashMap<>(),
1,
false,
+ workerStatus,
getNewReqeustId());
assertEquals(statusSystem.excludedWorkers.size(), 1);
@@ -559,6 +563,7 @@ public class DefaultMetaSystemSuiteJ {
new HashMap<>(),
1,
false,
+ workerStatus,
getNewReqeustId());
assertEquals(statusSystem.excludedWorkers.size(), 2);
@@ -574,6 +579,7 @@ public class DefaultMetaSystemSuiteJ {
new HashMap<>(),
1,
false,
+ workerStatus,
getNewReqeustId());
assertEquals(statusSystem.excludedWorkers.size(), 2);
@@ -589,6 +595,7 @@ public class DefaultMetaSystemSuiteJ {
new HashMap<>(),
1,
true,
+ workerStatus,
getNewReqeustId());
assertEquals(statusSystem.excludedWorkers.size(), 3);
diff --git
a/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/RatisMasterStatusSystemSuiteJ.java
b/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/RatisMasterStatusSystemSuiteJ.java
index c4fb8e29a..f192e91f9 100644
---
a/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/RatisMasterStatusSystemSuiteJ.java
+++
b/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/RatisMasterStatusSystemSuiteJ.java
@@ -34,6 +34,7 @@ import
org.apache.celeborn.common.exception.CelebornRuntimeException;
import org.apache.celeborn.common.identity.UserIdentifier;
import org.apache.celeborn.common.meta.DiskInfo;
import org.apache.celeborn.common.meta.WorkerInfo;
+import org.apache.celeborn.common.meta.WorkerStatus;
import org.apache.celeborn.common.quota.ResourceConsumption;
import org.apache.celeborn.common.rpc.RpcEndpointAddress;
import org.apache.celeborn.common.rpc.RpcEndpointRef;
@@ -211,6 +212,7 @@ public class RatisMasterStatusSystemSuiteJ {
private static final String APPID1 = "appId1";
private static final int SHUFFLEID1 = 1;
private static final String SHUFFLEKEY1 = APPID1 + "-" + SHUFFLEID1;
+ private static final WorkerStatus workerStatus =
WorkerStatus.normalWorkerStatus();
private String getNewReqeustId() {
return MasterClient.encodeRequestId(UUID.randomUUID().toString(),
callerId.incrementAndGet());
@@ -800,6 +802,7 @@ public class RatisMasterStatusSystemSuiteJ {
new HashMap<>(),
1,
false,
+ workerStatus,
getNewReqeustId());
Thread.sleep(3000L);
@@ -818,6 +821,7 @@ public class RatisMasterStatusSystemSuiteJ {
new HashMap<>(),
1,
false,
+ workerStatus,
getNewReqeustId());
Thread.sleep(3000L);
@@ -837,6 +841,7 @@ public class RatisMasterStatusSystemSuiteJ {
new HashMap<>(),
1,
false,
+ workerStatus,
getNewReqeustId());
Thread.sleep(3000L);
@@ -856,6 +861,7 @@ public class RatisMasterStatusSystemSuiteJ {
new HashMap<>(),
1,
true,
+ workerStatus,
getNewReqeustId());
Thread.sleep(3000L);
Assert.assertEquals(2, statusSystem.excludedWorkers.size());
diff --git
a/service/src/main/scala/org/apache/celeborn/server/common/HttpService.scala
b/service/src/main/scala/org/apache/celeborn/server/common/HttpService.scala
index 9ca1be8f7..31e931229 100644
--- a/service/src/main/scala/org/apache/celeborn/server/common/HttpService.scala
+++ b/service/src/main/scala/org/apache/celeborn/server/common/HttpService.scala
@@ -69,6 +69,11 @@ abstract class HttpService extends Service with Logging {
def exit(exitType: String): String = throw new
UnsupportedOperationException()
+ def handleWorkerEvent(workerEventType: String, workers: String): String =
+ throw new UnsupportedOperationException()
+
+ def getWorkerEventInfo(): String = throw new UnsupportedOperationException()
+
def startHttpServer(): Unit = {
val handlers =
if (metricsSystem.running) {
diff --git
a/service/src/main/scala/org/apache/celeborn/server/common/http/HttpEndpoint.scala
b/service/src/main/scala/org/apache/celeborn/server/common/http/HttpEndpoint.scala
index 1e4a48a63..29c1c1642 100644
---
a/service/src/main/scala/org/apache/celeborn/server/common/http/HttpEndpoint.scala
+++
b/service/src/main/scala/org/apache/celeborn/server/common/http/HttpEndpoint.scala
@@ -230,3 +230,23 @@ case object Exit extends HttpEndpoint {
override def handle(service: HttpService, parameters: Map[String, String]):
String =
service.exit(parameters.getOrElse("TYPE", ""))
}
+
+case object SendWorkerEvent extends HttpEndpoint {
+ override def path: String = "/sendWorkerEvent"
+
+ override def description(service: String): String =
+ "For Master(Leader) can send worker event to manager workers. Legal types
are 'None', 'Immediately', 'Decommission', 'DecommissionThenIdle', 'Graceful',
'Recommission'"
+
+ override def handle(service: HttpService, parameters: Map[String, String]):
String =
+ service.handleWorkerEvent(parameters.getOrElse("TYPE", ""),
parameters.getOrElse("WORKERS", ""))
+}
+
+case object WorkerEventInfo extends HttpEndpoint {
+ override def path: String = "/workerEventInfo"
+
+ override def description(service: String): String =
+ "List all worker event infos of the master."
+
+ override def handle(service: HttpService, parameters: Map[String, String]):
String =
+ service.getWorkerEventInfo()
+}
diff --git
a/service/src/main/scala/org/apache/celeborn/server/common/http/HttpUtils.scala
b/service/src/main/scala/org/apache/celeborn/server/common/http/HttpUtils.scala
index 534607d85..866150622 100644
---
a/service/src/main/scala/org/apache/celeborn/server/common/http/HttpUtils.scala
+++
b/service/src/main/scala/org/apache/celeborn/server/common/http/HttpUtils.scala
@@ -32,6 +32,8 @@ object HttpUtils {
ExcludedWorkers,
ShutdownWorkers,
Hostnames,
+ SendWorkerEvent,
+ WorkerEventInfo,
Exclude) ++ baseEndpoints
private val workerEndpoints: List[HttpEndpoint] =
List(
diff --git
a/service/src/test/scala/org/apache/celeborn/server/common/http/HttpUtilsSuite.scala
b/service/src/test/scala/org/apache/celeborn/server/common/http/HttpUtilsSuite.scala
index 9f025ae9a..182e80f58 100644
---
a/service/src/test/scala/org/apache/celeborn/server/common/http/HttpUtilsSuite.scala
+++
b/service/src/test/scala/org/apache/celeborn/server/common/http/HttpUtilsSuite.scala
@@ -70,9 +70,11 @@ class HttpUtilsSuite extends AnyFunSuite with Logging {
|/listTopDiskUsedApps List the top disk usage application ids. It
will return the top disk usage application ids for the cluster.
|/lostWorkers List all lost workers of the master.
|/masterGroupInfo List master group information of the service.
It will list all master's LEADER, FOLLOWER information.
+ |/sendWorkerEvent For Master(Leader) can send worker event to
manager workers. Legal types are 'None', 'Immediately', 'Decommission',
'DecommissionThenIdle', 'Graceful', 'Recommission'
|/shuffles List all running shuffle keys of the service.
It will return all running shuffle's key of the cluster.
|/shutdownWorkers List all shutdown workers of the master.
|/threadDump List the current thread dump of the master.
+ |/workerEventInfo List all worker event infos of the master.
|/workerInfo List worker information of the service. It will
list all registered workers 's information.
|""".stripMargin)
assert(HttpUtils.help(Service.WORKER) ==
@@ -91,4 +93,11 @@ class HttpUtilsSuite extends AnyFunSuite with Logging {
|/workerInfo List the worker information of the worker.
|""".stripMargin)
}
+
+ test("CELEBORN-1245: Support Master manage workers") {
+ checkParseUri(
+
"/sendWorkerEvent?type=decommission&workers=localhost:1001:1002:1003:1004",
+ "/sendWorkerEvent",
+ Map("TYPE" -> "decommission", "WORKERS" ->
"localhost:1001:1002:1003:1004"))
+ }
}
diff --git
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Worker.scala
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Worker.scala
index cca324445..ec170e64a 100644
---
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Worker.scala
+++
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/Worker.scala
@@ -39,7 +39,8 @@ import org.apache.celeborn.common.meta.{DiskInfo, WorkerInfo,
WorkerPartitionLoc
import org.apache.celeborn.common.metrics.MetricsSystem
import org.apache.celeborn.common.metrics.source.{JVMCPUSource, JVMSource,
ResourceConsumptionSource, SystemMiscSource, ThreadPoolSource}
import org.apache.celeborn.common.network.TransportContext
-import org.apache.celeborn.common.protocol.{PartitionType,
PbRegisterWorkerResponse, PbWorkerLostResponse, RpcNameConstants,
TransportModuleConstants}
+import org.apache.celeborn.common.protocol.{PartitionType,
PbRegisterWorkerResponse, PbWorkerLostResponse, RpcNameConstants,
TransportModuleConstants, WorkerEventType}
+import org.apache.celeborn.common.protocol.PbWorkerStatus.State
import org.apache.celeborn.common.protocol.message.ControlMessages._
import org.apache.celeborn.common.quota.ResourceConsumption
import org.apache.celeborn.common.rpc._
@@ -76,6 +77,7 @@ private[celeborn] class Worker(
metricsSystem.registerSource(new JVMCPUSource(conf,
MetricsSystem.ROLE_WORKER))
metricsSystem.registerSource(new SystemMiscSource(conf,
MetricsSystem.ROLE_WORKER))
+ val workerStatusManager = new WorkerStatusManager(conf)
val rpcEnv: RpcEnv = RpcEnv.create(
RpcNameConstants.WORKER_SYS,
workerArgs.host,
@@ -91,7 +93,6 @@ private[celeborn] class Worker(
private val WORKER_SHUTDOWN_PRIORITY = 100
val shutdown = new AtomicBoolean(false)
private val gracefulShutdown = conf.workerGracefulShutdown
- private var exitKind = CelebornExitKind.EXIT_IMMEDIATELY
if (gracefulShutdown) {
val checkPortMap = Map(
WORKER_RPC_PORT -> conf.workerRpcPort,
@@ -102,7 +103,6 @@ private[celeborn] class Worker(
!checkPortMap.values.exists(_ == 0),
"If enable graceful shutdown, the worker should use non-zero port. " +
s"${checkPortMap.map { case (k, v) => k.key + "=" + v }.mkString(",
")}")
- exitKind = CelebornExitKind.WORKER_GRACEFUL_SHUTDOWN
try {
val recoverRoot = new File(conf.workerGracefulShutdownRecoverPath)
if (!recoverRoot.exists()) {
@@ -350,6 +350,7 @@ private[celeborn] class Worker(
workerInfo.updateThenGetDiskInfos(storageManager.disksSnapshot().map {
disk =>
disk.mountPoint -> disk
}.toMap.asJava).values().asScala.toSeq ++ storageManager.hdfsDiskInfo
+ workerStatusManager.checkIfNeedTransitionStatus()
val response = masterClient.askSync[HeartbeatFromWorkerResponse](
HeartbeatFromWorker(
host,
@@ -361,10 +362,14 @@ private[celeborn] class Worker(
handleResourceConsumption(),
activeShuffleKeys,
estimatedAppDiskUsage,
- highWorkload),
+ highWorkload,
+ workerStatusManager.currentWorkerStatus),
classOf[HeartbeatFromWorkerResponse])
response.expiredShuffleKeys.asScala.foreach(shuffleKey =>
workerInfo.releaseSlots(shuffleKey))
cleanTaskQueue.put(response.expiredShuffleKeys)
+
+ val workerEvent = response.workerEvent
+ workerStatusManager.doTransition(workerEvent)
if (!response.registered) {
logError("Worker not registered in master, clean expired shuffle data
and register again.")
try {
@@ -425,6 +430,7 @@ private[celeborn] class Worker(
replicateHandler.init(this)
fetchHandler.init(this)
controller.init(this)
+ workerStatusManager.init(this)
logInfo("Worker started.")
rpcEnv.awaitTermination()
@@ -644,23 +650,17 @@ private[celeborn] class Worker(
override def exit(exitType: String): String = {
exitType.toUpperCase(Locale.ROOT) match {
case "DECOMMISSION" =>
- exitKind = CelebornExitKind.WORKER_DECOMMISSION
ShutdownHookManager.get().updateTimeout(
conf.workerDecommissionForceExitTimeout,
TimeUnit.MILLISECONDS)
+ workerStatusManager.doTransition(WorkerEventType.Decommission)
case "GRACEFUL" =>
- exitKind = CelebornExitKind.WORKER_GRACEFUL_SHUTDOWN
+ workerStatusManager.doTransition(WorkerEventType.Graceful)
case "IMMEDIATELY" =>
- exitKind = CelebornExitKind.EXIT_IMMEDIATELY
- case _ => // Use origin code
+ workerStatusManager.doTransition(WorkerEventType.Immediately)
+ case _ =>
+ workerStatusManager.doTransition(workerStatusManager.exitEventType)
}
- // Use the original EXIT_CODE
- new Thread() {
- override def run(): Unit = {
- Thread.sleep(10000)
- System.exit(0)
- }
- }.start()
val sb = new StringBuilder
sb.append("============================ Exit Worker
=============================\n")
sb.append(s"Exit worker by $exitType triggered: \n")
@@ -672,6 +672,8 @@ private[celeborn] class Worker(
// During shutdown, to avoid allocate slots in this worker,
// add this worker to master's excluded list. When restart, register
worker will
// make master remove this worker from excluded list.
+ logInfo("Worker start to shutdown gracefully")
+ workerStatusManager.transitionState(State.InGraceFul)
try {
masterClient.askSync(
ReportWorkerUnavailable(List(workerInfo).asJava),
@@ -701,9 +703,11 @@ private[celeborn] class Worker(
logWarning(s"Waiting for all PartitionLocation release cost
${waitTime}ms, " +
s"unreleased PartitionLocation: \n$partitionLocationInfo")
}
+
+ workerStatusManager.transitionState(State.Exit)
}
- def decommissionWorker(): Unit = {
+ def sendWorkerUnavailableToMaster(): Unit = {
try {
masterClient.askSync(
ReportWorkerUnavailable(List(workerInfo).asJava),
@@ -715,6 +719,12 @@ private[celeborn] class Worker(
s"\n${storageManager.shuffleKeySet().asScala.mkString("[", ", ",
"]")}",
e)
}
+ }
+
+ def decommissionWorker(): Unit = {
+ logInfo("Worker start to decommission")
+ workerStatusManager.transitionState(State.InDecommission)
+ sendWorkerUnavailableToMaster()
shutdown.set(true)
val interval = conf.workerDecommissionCheckInterval
val timeout = conf.workerDecommissionForceExitTimeout
@@ -732,12 +742,15 @@ private[celeborn] class Worker(
logWarning(s"Waiting for all shuffle expired cost ${waitTime}ms, " +
s"unreleased shuffle:
\n${storageManager.shuffleKeySet().asScala.mkString("[", ", ", "]")}")
}
+ workerStatusManager.transitionState(State.Exit)
}
def exitImmediately(): Unit = {
// During shutdown, to avoid allocate slots in this worker,
// add this worker to master's excluded list. When restart, register
worker will
// make master remove this worker from excluded list.
+ logInfo("Worker start to exit immediately")
+ workerStatusManager.transitionState(State.InExit)
try {
masterClient.askSync[PbWorkerLostResponse](
WorkerLost(
@@ -755,24 +768,27 @@ private[celeborn] class Worker(
e)
}
shutdown.set(true)
+ workerStatusManager.transitionState(State.Exit)
}
ShutdownHookManager.get().addShutdownHook(
new Thread(new Runnable {
override def run(): Unit = {
logInfo("Shutdown hook called.")
- exitKind match {
- case CelebornExitKind.WORKER_GRACEFUL_SHUTDOWN =>
- logInfo("Worker start to shutdown gracefully")
+ workerStatusManager.exitEventType match {
+ case WorkerEventType.Graceful =>
shutdownGracefully()
- case CelebornExitKind.WORKER_DECOMMISSION =>
- logInfo("Worker start to decommission")
+ case WorkerEventType.Decommission =>
decommissionWorker()
case _ =>
- logInfo("Worker start to exit immediately")
exitImmediately()
}
- stop(exitKind)
+
+ if (workerStatusManager.exitEventType == WorkerEventType.Graceful) {
+ stop(CelebornExitKind.WORKER_GRACEFUL_SHUTDOWN)
+ } else {
+ stop(CelebornExitKind.EXIT_IMMEDIATELY)
+ }
}
}),
WORKER_SHUTDOWN_PRIORITY)
diff --git
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/WorkerStatusManager.scala
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/WorkerStatusManager.scala
new file mode 100644
index 000000000..d6648c2fa
--- /dev/null
+++
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/WorkerStatusManager.scala
@@ -0,0 +1,158 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.celeborn.service.deploy.worker
+
+import java.util
+import java.util.concurrent.atomic.AtomicBoolean
+
+import scala.collection.immutable.HashSet
+
+import com.google.common.collect.Sets
+
+import org.apache.celeborn.common.CelebornConf
+import org.apache.celeborn.common.internal.Logging
+import org.apache.celeborn.common.meta.WorkerStatus
+import org.apache.celeborn.common.protocol.{PbWorkerStatus, WorkerEventType}
+import org.apache.celeborn.common.protocol.PbWorkerStatus.State
+import org.apache.celeborn.service.deploy.worker.storage.StorageManager
+
+private[celeborn] class WorkerStatusManager(conf: CelebornConf) extends
Logging {
+
+ var currentWorkerStatus = WorkerStatus.normalWorkerStatus()
+ var exitEventType = WorkerEventType.Immediately
+ private var worker: Worker = _
+ private var shutdown: AtomicBoolean = _
+ private var storageManager: StorageManager = _
+ private val gracefulShutdown = conf.workerGracefulShutdown
+ if (gracefulShutdown) {
+ exitEventType = WorkerEventType.Graceful
+ }
+
+ private val transitionStateMap = new util.HashMap[State,
util.HashSet[State]]()
+ transitionStateMap.put(
+ State.Normal,
+ Sets.newHashSet(
+ State.InGraceFul,
+ State.InExit,
+ State.InDecommission,
+ State.InDecommissionThenIdle))
+ transitionStateMap.put(
+ State.InDecommissionThenIdle,
+ Sets.newHashSet(State.Normal, State.Idle, State.InGraceFul, State.InExit,
State.InDecommission))
+ transitionStateMap.put(
+ State.Idle,
+ Sets.newHashSet(State.Normal, State.InGraceFul, State.InExit,
State.InDecommission))
+ transitionStateMap.put(State.InGraceFul, Sets.newHashSet(State.Exit))
+ transitionStateMap.put(State.InDecommission, Sets.newHashSet(State.Exit))
+ transitionStateMap.put(State.InExit, Sets.newHashSet(State.Exit))
+
+ private val exitStatus = HashSet(
+ State.InExit,
+ State.Exit,
+ State.InDecommission,
+ State.InGraceFul)
+
+ def init(worker: Worker): Unit = {
+ this.worker = worker
+ shutdown = worker.shutdown
+ storageManager = worker.storageManager
+ }
+
+ def doTransition(eventType: WorkerEventType): Unit = this.synchronized {
+ if (inExitStatus()) {
+ logDebug(s"Worker receive event: $eventType, but in exit State:
${getWorkerState()} ")
+ } else {
+ logDebug(s"Worker receive event: $eventType, currentState:
${getWorkerState()} ")
+ checkIfNeedTransitionStatus()
+ val currentState = getWorkerState()
+ eventType match {
+ case WorkerEventType.DecommissionThenIdle if currentState ==
State.Normal =>
+ decommissionWorkerThenIdle()
+ case WorkerEventType.Recommission
+ if currentState == State.InDecommissionThenIdle || currentState ==
State.Idle =>
+ recommissionWorker()
+ case WorkerEventType.Graceful | WorkerEventType.Immediately |
WorkerEventType.Decommission =>
+ exit(eventType)
+ case _ =>
+ logDebug(s"Worker receive event: $eventType, and has nothing to do ")
+ }
+ }
+ }
+
+ def checkIfNeedTransitionStatus(): Unit = this.synchronized {
+ currentWorkerStatus.getState match {
+ case State.InDecommissionThenIdle if decommissionThenIdleFinished() =>
+ transitionState(State.Idle)
+ case _ =>
+ }
+ }
+
+ private def exit(eventType: WorkerEventType): Unit = {
+ exitEventType = eventType
+ exitEventType match {
+ case WorkerEventType.Immediately => transitionState(State.InExit)
+ case WorkerEventType.Graceful => transitionState(State.InGraceFul)
+ case WorkerEventType.Decommission =>
transitionState(State.InDecommission)
+ case _ => // ignore
+ }
+
+ // Compatible with current exit logic
+ // trigger shutdown hook to exit
+ new Thread() {
+ override def run(): Unit = {
+ Thread.sleep(10000)
+ System.exit(0)
+ }
+ }.start()
+ }
+
+ def transitionState(state: State): Unit = this.synchronized {
+ val allowStates = transitionStateMap.get(currentWorkerStatus.getState)
+ if (allowStates != null && allowStates.contains(state)) {
+ logInfo(s"Worker transition status from ${currentWorkerStatus.getState}
to $state.")
+ currentWorkerStatus = new WorkerStatus(state.getNumber,
System.currentTimeMillis())
+ } else {
+ logWarning(
+ s"Worker transition status from ${currentWorkerStatus.getState} to
$state is not allowed.")
+ }
+ }
+
+ private def recommissionWorker(): Unit = this.synchronized {
+ shutdown.set(false)
+ transitionState(State.Normal)
+ }
+
+ private def decommissionWorkerThenIdle(): Unit = this.synchronized {
+ shutdown.set(true)
+ transitionState(State.InDecommissionThenIdle)
+ worker.sendWorkerUnavailableToMaster()
+ checkIfNeedTransitionStatus()
+ }
+
+ private def decommissionThenIdleFinished(): Boolean = this.synchronized {
+ shutdown.get() && (storageManager.shuffleKeySet().isEmpty ||
currentWorkerStatus.getState == State.Idle)
+ }
+
+ def getWorkerState(): State = {
+ currentWorkerStatus.getState
+ }
+
+ private def inExitStatus(): Boolean = {
+ exitStatus.contains(currentWorkerStatus.getState)
+ }
+}
diff --git
a/worker/src/test/scala/org/apache/celeborn/service/deploy/worker/WorkerStatusManagerSuite.scala
b/worker/src/test/scala/org/apache/celeborn/service/deploy/worker/WorkerStatusManagerSuite.scala
new file mode 100644
index 000000000..b1e7d1fc6
--- /dev/null
+++
b/worker/src/test/scala/org/apache/celeborn/service/deploy/worker/WorkerStatusManagerSuite.scala
@@ -0,0 +1,64 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.celeborn.service.deploy.worker
+
+import java.util.concurrent.atomic.AtomicBoolean
+
+import com.google.common.collect.Sets
+import org.junit.Assert
+import org.mockito.MockitoSugar._
+import org.scalatest.funsuite.AnyFunSuite
+
+import org.apache.celeborn.common.CelebornConf
+import org.apache.celeborn.common.protocol.{PbWorkerStatus, WorkerEventType}
+import org.apache.celeborn.service.deploy.worker.storage.StorageManager
+
+class WorkerStatusManagerSuite extends AnyFunSuite {
+ private var worker: Worker = _
+ private val conf = new CelebornConf()
+
+ test("Test Worker status transition to none exit status") {
+ worker = mock[Worker]
+ val storageManager = mock[StorageManager]
+ val shuffleKeys = Sets.newHashSet("test")
+ when(storageManager.shuffleKeySet()).thenReturn(shuffleKeys)
+ when(worker.storageManager).thenReturn(storageManager)
+ when(worker.shutdown).thenReturn(new AtomicBoolean())
+
+ val statusManager = new WorkerStatusManager(conf)
+ statusManager.init(worker)
+
+ statusManager.doTransition(WorkerEventType.DecommissionThenIdle)
+ Assert.assertEquals(statusManager.getWorkerState(),
PbWorkerStatus.State.InDecommissionThenIdle)
+
+ // Rerun state Transition
+ statusManager.doTransition(WorkerEventType.DecommissionThenIdle)
+ Assert.assertEquals(statusManager.getWorkerState(),
PbWorkerStatus.State.InDecommissionThenIdle)
+
+ // Reset shuffleKeys
+ shuffleKeys.clear()
+ statusManager.doTransition(WorkerEventType.DecommissionThenIdle)
+ Assert.assertEquals(statusManager.getWorkerState(),
PbWorkerStatus.State.Idle)
+
+ statusManager.doTransition(WorkerEventType.Recommission)
+ Assert.assertEquals(statusManager.getWorkerState(),
PbWorkerStatus.State.Normal)
+
+ statusManager.doTransition(WorkerEventType.Recommission)
+ Assert.assertEquals(statusManager.getWorkerState(),
PbWorkerStatus.State.Normal)
+ }
+}