Skip to content

Commit 99d7186

Browse files
committed
[CELEBORN-2349] Support worker registered tags during startup
1 parent 3820244 commit 99d7186

15 files changed

Lines changed: 180 additions & 178 deletions

File tree

common/src/main/proto/TransportMessages.proto

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -189,6 +189,7 @@ message PbWorkerInfo {
189189
int32 internalPort = 8;
190190
string networkLocation = 9;
191191
int64 nextInterruptionNotice = 10; // Unix timestamp when disruption is expected to be initiated
192+
repeated string tags = 11;
192193
}
193194

194195
message PbFileGroup {
@@ -207,6 +208,7 @@ message PbRegisterWorker {
207208
map<string, PbResourceConsumption> userResourceConsumption = 8;
208209
int32 internalPort = 10;
209210
string networkLocation = 11;
211+
repeated string tags = 12;
210212
}
211213

212214
message PbMetaRegisterWorkerRequest {
@@ -219,6 +221,7 @@ message PbMetaRegisterWorkerRequest {
219221
map<string, PbResourceConsumption> userResourceConsumption = 7;
220222
int32 internalPort = 8;
221223
string networkLocation = 9;
224+
repeated string tags = 12;
222225
}
223226

224227
message PbHeartbeatFromWorker {

common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1527,6 +1527,8 @@ class CelebornConf(loadDefaults: Boolean) extends Cloneable with Logging with Se
15271527
def clientInputStreamCreationWindow = get(CLIENT_INPUTSTREAM_CREATION_WINDOW)
15281528

15291529
def tagsEnabled: Boolean = get(TAGS_ENABLED)
1530+
def tagsWorkerRegistrationEnabled: Boolean = get(TAGS_WORKER_REGISTRATION_ENABLED)
1531+
def workerTags: Seq[String] = get(WORKER_TAGS)
15301532
def tagsExpr: String = get(TAGS_EXPR)
15311533
def preferClientTagsExpr: Boolean = get(PREFER_CLIENT_TAGS_EXPR)
15321534

@@ -6901,6 +6903,24 @@ object CelebornConf extends Logging {
69016903
.booleanConf
69026904
.createWithDefault(true)
69036905

6906+
val TAGS_WORKER_REGISTRATION_ENABLED: ConfigEntry[Boolean] =
6907+
buildConf("celeborn.tags.worker.registration.enabled")
6908+
.categories("master")
6909+
.version("0.7.0")
6910+
.doc("When true, the master honors tags advertised by workers at registration " +
6911+
"(merged with the config-store tags). When false, worker-supplied tags are ignored.")
6912+
.booleanConf
6913+
.createWithDefault(true)
6914+
6915+
val WORKER_TAGS: ConfigEntry[Seq[String]] =
6916+
buildConf("celeborn.worker.tags")
6917+
.categories("worker")
6918+
.version("0.7.0")
6919+
.doc("Comma-separated tags this worker supplies to the master at registration.")
6920+
.stringConf
6921+
.toSequence
6922+
.createWithDefault(Seq.empty)
6923+
69046924
val TAGS_EXPR: ConfigEntry[String] =
69056925
buildConf("celeborn.tags.tagsExpr")
69066926
.categories("master", "client")

common/src/main/scala/org/apache/celeborn/common/meta/WorkerInfo.scala

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -47,7 +47,8 @@ class WorkerInfo(
4747
var nextInterruptionNotice = Long.MaxValue
4848
var lastHeartbeat: Long = 0
4949
var workerStatus = WorkerStatus.normalWorkerStatus()
50-
var isHighWorkLoad: Boolean = false;
50+
var isHighWorkLoad: Boolean = false
51+
var tags: util.Set[String] = new util.HashSet[String]()
5152
val diskInfos = {
5253
if (_diskInfos != null) JavaUtils.newConcurrentHashMap[String, DiskInfo](_diskInfos)
5354
else null

common/src/main/scala/org/apache/celeborn/common/protocol/message/ControlMessages.scala

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -93,6 +93,7 @@ object ControlMessages extends Logging {
9393
networkLocation: String,
9494
disks: Map[String, DiskInfo],
9595
userResourceConsumption: Map[UserIdentifier, ResourceConsumption],
96+
tags: Set[String],
9697
requestId: String): PbRegisterWorker = {
9798
val pbDisks = disks.values.map(PbSerDeUtils.toPbDiskInfo).asJava
9899
val pbUserResourceConsumption =
@@ -107,6 +108,7 @@ object ControlMessages extends Logging {
107108
.setNetworkLocation(networkLocation)
108109
.addAllDisks(pbDisks)
109110
.putAllUserResourceConsumption(pbUserResourceConsumption)
111+
.addAllTags(tags.asJava)
110112
.setRequestId(requestId)
111113
.build()
112114
}

common/src/main/scala/org/apache/celeborn/common/util/PbSerDeUtils.scala

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -296,6 +296,7 @@ object PbSerDeUtils {
296296
} else {
297297
workerInfo.nextInterruptionNotice = pbWorkerInfo.getNextInterruptionNotice
298298
}
299+
workerInfo.tags = new util.HashSet[String](pbWorkerInfo.getTagsList)
299300
workerInfo
300301
}
301302

@@ -310,6 +311,7 @@ object PbSerDeUtils {
310311
.setPushPort(workerInfo.pushPort)
311312
.setReplicatePort(workerInfo.replicatePort)
312313
.setInternalPort(workerInfo.internalPort)
314+
.addAllTags(workerInfo.tags)
313315
if (masterPersistWorkerNetworkLocation) {
314316
builder.setNetworkLocation(workerInfo.networkLocation)
315317
}

master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/AbstractMetaManager.java

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -357,7 +357,8 @@ public void updateRegisterWorkerMeta(
357357
int replicatePort,
358358
int internalPort,
359359
String networkLocation,
360-
Map<String, DiskInfo> disks) {
360+
Map<String, DiskInfo> disks,
361+
Set<String> tags) {
361362
WorkerInfo workerInfo =
362363
new WorkerInfo(
363364
host,
@@ -369,6 +370,7 @@ public void updateRegisterWorkerMeta(
369370
disks,
370371
new HashMap<>());
371372
workerInfo.lastHeartbeat_$eq(System.currentTimeMillis());
373+
workerInfo.tags_$eq(new HashSet<>(tags));
372374
if (networkLocation != null
373375
&& !networkLocation.isEmpty()
374376
&& !NetworkTopology.DEFAULT_RACK.equals(networkLocation)) {

master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/IMetadataHandler.java

Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,8 +17,10 @@
1717

1818
package org.apache.celeborn.service.deploy.master.clustermeta;
1919

20+
import java.util.Collections;
2021
import java.util.List;
2122
import java.util.Map;
23+
import java.util.Set;
2224

2325
import org.apache.celeborn.common.identity.UserIdentifier;
2426
import org.apache.celeborn.common.meta.ApplicationMeta;
@@ -77,6 +79,31 @@ void handleWorkerHeartbeat(
7779
WorkerStatus workerStatus,
7880
String requestId);
7981

82+
default void handleRegisterWorker(
83+
String host,
84+
int rpcPort,
85+
int pushPort,
86+
int fetchPort,
87+
int replicatePort,
88+
int internalPort,
89+
String networkLocation,
90+
Map<String, DiskInfo> disks,
91+
Map<UserIdentifier, ResourceConsumption> userResourceConsumption,
92+
String requestId) {
93+
handleRegisterWorker(
94+
host,
95+
rpcPort,
96+
pushPort,
97+
fetchPort,
98+
replicatePort,
99+
internalPort,
100+
networkLocation,
101+
disks,
102+
userResourceConsumption,
103+
Collections.emptySet(),
104+
requestId);
105+
}
106+
80107
void handleRegisterWorker(
81108
String host,
82109
int rpcPort,
@@ -87,6 +114,7 @@ void handleRegisterWorker(
87114
String networkLocation,
88115
Map<String, DiskInfo> disks,
89116
Map<UserIdentifier, ResourceConsumption> userResourceConsumption,
117+
Set<String> tags,
90118
String requestId);
91119

92120
void handleReportWorkerUnavailable(List<WorkerInfo> failedNodes, String requestId);

master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/SingleMasterMetaManager.java

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919

2020
import java.util.List;
2121
import java.util.Map;
22+
import java.util.Set;
2223

2324
import org.slf4j.Logger;
2425
import org.slf4j.LoggerFactory;
@@ -162,9 +163,18 @@ public void handleRegisterWorker(
162163
String networkLocation,
163164
Map<String, DiskInfo> disks,
164165
Map<UserIdentifier, ResourceConsumption> userResourceConsumption,
166+
Set<String> tags,
165167
String requestId) {
166168
updateRegisterWorkerMeta(
167-
host, rpcPort, pushPort, fetchPort, replicatePort, internalPort, networkLocation, disks);
169+
host,
170+
rpcPort,
171+
pushPort,
172+
fetchPort,
173+
replicatePort,
174+
internalPort,
175+
networkLocation,
176+
disks,
177+
tags);
168178
updateWorkerResourceConsumptions(
169179
host, rpcPort, pushPort, fetchPort, replicatePort, userResourceConsumption);
170180
}

master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/HAMasterMetaManager.java

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919

2020
import java.util.List;
2121
import java.util.Map;
22+
import java.util.Set;
2223
import java.util.stream.Collectors;
2324

2425
import org.slf4j.Logger;
@@ -349,6 +350,7 @@ public void handleRegisterWorker(
349350
String networkLocation,
350351
Map<String, DiskInfo> disks,
351352
Map<UserIdentifier, ResourceConsumption> userResourceConsumption,
353+
Set<String> tags,
352354
String requestId) {
353355
try {
354356
ratisServer.submitRequest(
@@ -365,6 +367,7 @@ public void handleRegisterWorker(
365367
.setInternalPort(internalPort)
366368
.setNetworkLocation(networkLocation)
367369
.putAllDisks(MetaUtil.toPbDiskInfos(disks))
370+
.addAllTags(tags)
368371
.build())
369372
.build());
370373
updateWorkerResourceConsumptions(

master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/MetaHandler.java

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -258,6 +258,7 @@ public org.apache.celeborn.common.protocol.PbMetaRequestResponse handleWriteRequ
258258
int internalPort = request.getRegisterWorkerRequest().getInternalPort();
259259
Map<String, PbDiskInfo> pbDiskInfo = request.getRegisterWorkerRequest().getDisksMap();
260260
diskInfos = MetaUtil.fromPbDiskInfoMap(pbDiskInfo);
261+
Set<String> tags = new HashSet<>(request.getRegisterWorkerRequest().getTagsList());
261262
LOG.debug(
262263
"Handle worker register for {} {} {} {} {} {} {}",
263264
host,
@@ -275,7 +276,8 @@ public org.apache.celeborn.common.protocol.PbMetaRequestResponse handleWriteRequ
275276
replicatePort,
276277
internalPort,
277278
networkLocation,
278-
diskInfos);
279+
diskInfos,
280+
tags);
279281
break;
280282

281283
case ReportWorkerUnavailable:

0 commit comments

Comments
 (0)