From 304a161521e5102c60b4fdb1367b78bcd8de22f8 Mon Sep 17 00:00:00 2001 From: yew1eb Date: Thu, 20 Aug 2026 10:24:22 +0800 Subject: [PATCH] [CELEBORN-2429][MASTER] Batch worker/application heartbeats into aggregated raft log entries Aggregate heartbeats on the leader over a short time window (default 1s) into one BatchHeartbeat raft entry, cutting raft write volume by ~100x at peak. Off by default: celeborn.master.ha.heartbeat.batch.enabled. --- common/src/main/proto/TransportMessages.proto | 7 + .../apache/celeborn/common/CelebornConf.scala | 26 ++ docs/configuration/ha.md | 2 + .../clustermeta/ha/HAMasterMetaManager.java | 67 +++-- .../clustermeta/ha/HeartbeatAggregator.java | 152 ++++++++++++ .../master/clustermeta/ha/MetaHandler.java | 134 +++++----- master/src/main/proto/Resource.proto | 8 + .../service/deploy/master/Master.scala | 9 + .../ha/HeartbeatAggregatorSuiteJ.java | 230 ++++++++++++++++++ .../ha/MasterStateMachineSuiteJ.java | 110 +++++++++ .../clustermeta/ha/RatisBaseSuiteJ.java | 3 +- 11 files changed, 666 insertions(+), 82 deletions(-) create mode 100644 master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/HeartbeatAggregator.java create mode 100644 master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/HeartbeatAggregatorSuiteJ.java diff --git a/common/src/main/proto/TransportMessages.proto b/common/src/main/proto/TransportMessages.proto index a813a9e5015..9d37e3ca6a6 100644 --- a/common/src/main/proto/TransportMessages.proto +++ b/common/src/main/proto/TransportMessages.proto @@ -463,6 +463,11 @@ message PbMetaBatchUnregisterShuffles { repeated string shuffleKeys = 1; } +message PbBatchHeartbeatRequest { + repeated PbMetaWorkerHeartbeatRequest workerHeartbeats = 1; + repeated PbMetaAppHeartbeatRequest appHeartbeats = 2; +} + message PbUnregisterShuffleResponse { int32 status = 1; } @@ -998,6 +1003,7 @@ enum PbMetaRequestType { BatchUnRegisterShuffle = 28; ReviseLostShuffles = 29; RegisterApplicationInfo = 30; + BatchHeartbeat = 31; } message PbMetaRequest { @@ -1019,6 +1025,7 @@ message PbMetaRequest { PbReportWorkerDecommission reportWorkerDecommissionRequest = 24; PbMetaBatchUnregisterShuffles batchUnregisterShuffleRequest = 25; PbRegisterApplicationInfo registerApplicationInfoRequest = 26; + PbBatchHeartbeatRequest batchHeartbeatRequest = 27; PbReviseLostShuffles reviseLostShufflesRequest = 102; } diff --git a/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala b/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala index 6dce12fec67..5b3665c3e7d 100644 --- a/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala +++ b/common/src/main/scala/org/apache/celeborn/common/CelebornConf.scala @@ -815,6 +815,8 @@ class CelebornConf(loadDefaults: Boolean) extends Cloneable with Logging with Se def haMasterRatisSnapshotAutoTriggerThreshold: Long = get(HA_MASTER_RATIS_SNAPSHOT_AUTO_TRIGGER_THRESHOLD) def haMasterRatisSnapshotRetentionFileNum: Int = get(HA_MASTER_RATIS_SNAPSHOT_RETENTION_FILE_NUM) + def masterHaHeartbeatBatchEnabled: Boolean = get(MASTER_HA_HEARTBEAT_BATCH_ENABLED) + def masterHaHeartbeatBatchIntervalMs: Long = get(MASTER_HA_HEARTBEAT_BATCH_INTERVAL) def masterPersistWorkerNetworkLocation: Boolean = get(MASTER_PERSIST_WORKER_NETWORK_LOCATION) def haRatisCustomConfigs: JMap[String, String] = { @@ -3086,6 +3088,30 @@ object CelebornConf extends Logging { .intConf .createWithDefault(3) + val MASTER_HA_HEARTBEAT_BATCH_ENABLED: ConfigEntry[Boolean] = + buildConf("celeborn.master.ha.heartbeat.batch.enabled") + .categories("ha") + .version("1.0.0") + .doc( + "Whether to aggregate worker/app heartbeats on the raft leader and flush them " + + "periodically as a single BatchHeartbeat raft log entry, reducing raft log entries, " + + "fsyncs and state machine applies by roughly the number of heartbeats per flush window.") + .booleanConf + .createWithDefault(false) + + val MASTER_HA_HEARTBEAT_BATCH_INTERVAL: ConfigEntry[Long] = + buildConf("celeborn.master.ha.heartbeat.batch.interval") + .categories("ha") + .version("1.0.0") + .doc( + "The interval at which aggregated heartbeats are flushed as a single raft log entry. " + + "Heartbeat replies do not wait for raft replication when batching is enabled, so a " + + "larger interval only means the leader may lose at most one interval of heartbeat " + + "metadata on failover, which is negligible compared to the worker/app heartbeat " + + "timeouts (120s/300s by default).") + .timeConf(TimeUnit.MILLISECONDS) + .createWithDefaultString("1s") + val MASTER_PERSIST_WORKER_NETWORK_LOCATION: ConfigEntry[Boolean] = buildConf("celeborn.master.persist.workerNetworkLocation") .categories("master") diff --git a/docs/configuration/ha.md b/docs/configuration/ha.md index ed8ec7dbec8..3aef830a5b9 100644 --- a/docs/configuration/ha.md +++ b/docs/configuration/ha.md @@ -22,6 +22,8 @@ license: | | celeborn.master.ha.enabled | false | false | When true, master nodes run as Raft cluster mode. | 0.3.0 | celeborn.ha.enabled | | celeborn.master.ha.graceful.shutdown.enabled | false | false | When true, the master will transfer Raft leadership before shutting down gracefully. This reduces chances of client side failures by avoiding the Raft election window where no leader is available. | 0.7.0 | | | celeborn.master.ha.graceful.shutdown.timeout | 30s | false | Timeout for the master graceful shutdown process including Raft leadership transfer. Used as the shutdown hook timeout and the transfer-leadership request timeout. | 0.7.0 | | +| celeborn.master.ha.heartbeat.batch.enabled | false | false | Whether to aggregate worker/app heartbeats on the raft leader and flush them periodically as a single BatchHeartbeat raft log entry, reducing raft log entries, fsyncs and state machine applies by roughly the number of heartbeats per flush window. | 1.0.0 | | +| celeborn.master.ha.heartbeat.batch.interval | 1s | false | The interval at which aggregated heartbeats are flushed as a single raft log entry. Heartbeat replies do not wait for raft replication when batching is enabled, so a larger interval only means the leader may lose at most one interval of heartbeat metadata on failover, which is negligible compared to the worker/app heartbeat timeouts (120s/300s by default). | 1.0.0 | | | celeborn.master.ha.node.<id>.host | <required> | false | Host to bind of master node in HA mode. | 0.3.0 | celeborn.ha.master.node.<id>.host | | celeborn.master.ha.node.<id>.internal.port | 8097 | false | Internal port for the workers and other masters to bind to a master node in HA mode. | 0.5.0 | | | celeborn.master.ha.node.<id>.port | 9097 | false | Port to bind of master node in HA mode. | 0.3.0 | celeborn.ha.master.node.<id>.port | 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 3372143aa23..75f44aaebcd 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 @@ -46,6 +46,8 @@ public class HAMasterMetaManager extends AbstractMetaManager { protected HARaftServer ratisServer; + private HeartbeatAggregator heartbeatAggregator; + public HAMasterMetaManager(RpcEnv rpcEnv, CelebornConf conf) { this(rpcEnv, conf, new CelebornRackResolver(conf)); } @@ -68,6 +70,13 @@ public HARaftServer getRatisServer() { public void setRatisServer(HARaftServer ratisServer) { this.ratisServer = ratisServer; + if (conf.masterHaHeartbeatBatchEnabled()) { + this.heartbeatAggregator = new HeartbeatAggregator(ratisServer, conf); + } + } + + public HeartbeatAggregator getHeartbeatAggregator() { + return heartbeatAggregator; } @Override @@ -168,22 +177,27 @@ public void handleAppHeartbeat( Map applicationFallbackCounts, long time, String requestId) { + ResourceProtos.AppHeartbeatRequest appHeartbeatRequest = + ResourceProtos.AppHeartbeatRequest.newBuilder() + .setAppId(appId) + .setTime(time) + .setTotalWritten(totalWritten) + .setFileCount(fileCount) + .setShuffleCount(shuffleCount) + .setApplicationCount(applicationCount) + .putAllShuffleFallbackCounts(shuffleFallbackCounts) + .putAllApplicationFallbackCounts(applicationFallbackCounts) + .build(); + if (heartbeatAggregator != null) { + heartbeatAggregator.offerAppHeartbeat(appHeartbeatRequest); + return; + } try { ratisServer.submitRequest( ResourceRequest.newBuilder() .setCmdType(Type.AppHeartbeat) .setRequestId(requestId) - .setAppHeartbeatRequest( - ResourceProtos.AppHeartbeatRequest.newBuilder() - .setAppId(appId) - .setTime(time) - .setTotalWritten(totalWritten) - .setFileCount(fileCount) - .setShuffleCount(shuffleCount) - .setApplicationCount(applicationCount) - .putAllShuffleFallbackCounts(shuffleFallbackCounts) - .putAllApplicationFallbackCounts(applicationFallbackCounts) - .build()) + .setAppHeartbeatRequest(appHeartbeatRequest) .build()); } catch (CelebornRuntimeException e) { LOG.error("Handle heartbeat for {} failed!", appId, e); @@ -312,23 +326,30 @@ public void handleWorkerHeartbeat( boolean highWorkload, WorkerStatus workerStatus, String requestId) { + ResourceProtos.WorkerHeartbeatRequest workerHeartbeatRequest = + ResourceProtos.WorkerHeartbeatRequest.newBuilder() + .setHost(host) + .setRpcPort(rpcPort) + .setPushPort(pushPort) + .setFetchPort(fetchPort) + .setReplicatePort(replicatePort) + .putAllDisks(MetaUtil.toPbDiskInfos(disks)) + .setWorkerStatus(MetaUtil.toPbWorkerStatus(workerStatus)) + .setTime(time) + .setHighWorkload(highWorkload) + .build(); + if (heartbeatAggregator != null) { + heartbeatAggregator.offerWorkerHeartbeat(workerHeartbeatRequest); + updateWorkerResourceConsumptions( + host, rpcPort, pushPort, fetchPort, replicatePort, userResourceConsumption); + return; + } try { ratisServer.submitRequest( ResourceRequest.newBuilder() .setCmdType(Type.WorkerHeartbeat) .setRequestId(requestId) - .setWorkerHeartbeatRequest( - ResourceProtos.WorkerHeartbeatRequest.newBuilder() - .setHost(host) - .setRpcPort(rpcPort) - .setPushPort(pushPort) - .setFetchPort(fetchPort) - .setReplicatePort(replicatePort) - .putAllDisks(MetaUtil.toPbDiskInfos(disks)) - .setWorkerStatus(MetaUtil.toPbWorkerStatus(workerStatus)) - .setTime(time) - .setHighWorkload(highWorkload) - .build()) + .setWorkerHeartbeatRequest(workerHeartbeatRequest) .build()); updateWorkerResourceConsumptions( host, rpcPort, pushPort, fetchPort, replicatePort, userResourceConsumption); diff --git a/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/HeartbeatAggregator.java b/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/HeartbeatAggregator.java new file mode 100644 index 00000000000..97fb55461bd --- /dev/null +++ b/master/src/main/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/HeartbeatAggregator.java @@ -0,0 +1,152 @@ +/* + * 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.master.clustermeta.ha; + +import java.util.ArrayList; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.locks.ReentrantLock; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import org.apache.celeborn.common.CelebornConf; +import org.apache.celeborn.common.client.MasterClient; +import org.apache.celeborn.common.util.ThreadUtils; +import org.apache.celeborn.service.deploy.master.clustermeta.ResourceProtos; +import org.apache.celeborn.service.deploy.master.clustermeta.ResourceProtos.ResourceRequest; +import org.apache.celeborn.service.deploy.master.clustermeta.ResourceProtos.Type; + +public class HeartbeatAggregator { + private static final Logger LOG = LoggerFactory.getLogger(HeartbeatAggregator.class); + + private final HARaftServer ratisServer; + private final long batchIntervalMs; + + private final ReentrantLock pendingLock = new ReentrantLock(); + private Map workerHeartbeats = new HashMap<>(); + private Map appHeartbeats = new HashMap<>(); + + private final ScheduledExecutorService flushExecutor; + + public HeartbeatAggregator(HARaftServer ratisServer, CelebornConf conf) { + this.ratisServer = ratisServer; + this.batchIntervalMs = conf.masterHaHeartbeatBatchIntervalMs(); + this.flushExecutor = + ThreadUtils.newDaemonSingleThreadScheduledExecutor("master-heartbeat-aggregator"); + this.flushExecutor.scheduleWithFixedDelay( + this::flushSafely, batchIntervalMs, batchIntervalMs, TimeUnit.MILLISECONDS); + LOG.info("HeartbeatAggregator started, flush interval {} ms.", batchIntervalMs); + } + + public void offerWorkerHeartbeat(ResourceProtos.WorkerHeartbeatRequest heartbeat) { + pendingLock.lock(); + try { + workerHeartbeats.put(workerKey(heartbeat), heartbeat); + } finally { + pendingLock.unlock(); + } + } + + public void offerAppHeartbeat(ResourceProtos.AppHeartbeatRequest heartbeat) { + pendingLock.lock(); + try { + appHeartbeats.put(heartbeat.getAppId(), heartbeat); + } finally { + pendingLock.unlock(); + } + } + + private static String workerKey(ResourceProtos.WorkerHeartbeatRequest heartbeat) { + return heartbeat.getHost() + + ":" + + heartbeat.getRpcPort() + + ":" + + heartbeat.getPushPort() + + ":" + + heartbeat.getFetchPort() + + ":" + + heartbeat.getReplicatePort(); + } + + public void stop() { + flushExecutor.shutdownNow(); + } + + private void flushSafely() { + try { + flush(); + } catch (Throwable t) { + // Dropped batches self-heal next interval. + LOG.error("Failed to flush aggregated heartbeats, dropping this batch.", t); + } + } + + private void flush() { + List drainedWorkers; + List drainedApps; + pendingLock.lock(); + try { + if (workerHeartbeats.isEmpty() && appHeartbeats.isEmpty()) { + return; + } + drainedWorkers = new ArrayList<>(workerHeartbeats.values()); + workerHeartbeats.clear(); + drainedApps = new ArrayList<>(appHeartbeats.values()); + appHeartbeats.clear(); + } finally { + pendingLock.unlock(); + } + + if (!ratisServer.isLeader()) { + return; + } + + ResourceRequest batchRequest = + ResourceRequest.newBuilder() + .setCmdType(Type.BatchHeartbeat) + .setRequestId(MasterClient.genRequestId()) + .setBatchHeartbeatRequest( + ResourceProtos.BatchHeartbeatRequest.newBuilder() + .addAllWorkerHeartbeats(drainedWorkers) + .addAllAppHeartbeats(drainedApps) + .build()) + .build(); + long startNs = System.nanoTime(); + ratisServer.submitRequest(batchRequest); + long elapsedMs = TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - startNs); + if (elapsedMs > batchIntervalMs) { + LOG.warn( + "Submitting aggregated heartbeats ({} worker, {} app) took {} ms; raft commits are " + + "slower than the flush interval {} ms.", + drainedWorkers.size(), + drainedApps.size(), + elapsedMs, + batchIntervalMs); + } + if (LOG.isDebugEnabled()) { + LOG.debug( + "Flushed aggregated heartbeats, {} worker heartbeats, {} app heartbeats.", + drainedWorkers.size(), + drainedApps.size()); + } + } +} 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 86097117e88..26de66eda33 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 @@ -156,32 +156,24 @@ public org.apache.celeborn.common.protocol.PbMetaRequestResponse handleWriteRequ break; case AppHeartbeat: - appId = request.getAppHeartbeatRequest().getAppId(); - long time = request.getAppHeartbeatRequest().getTime(); - long totalWritten = request.getAppHeartbeatRequest().getTotalWritten(); - long fileCount = request.getAppHeartbeatRequest().getFileCount(); - long shuffleCount = request.getAppHeartbeatRequest().getShuffleCount(); - long applicationCount = request.getAppHeartbeatRequest().getApplicationCount(); - LOG.debug("Handle app heartbeat for {} with shuffle count {}", appId, shuffleCount); - Map shuffleFallbackCounts = - request.getAppHeartbeatRequest().getShuffleFallbackCountsMap(); - if (CollectionUtils.isNotEmpty(shuffleFallbackCounts)) { - LOG.warn( - "{} shuffle fallbacks in app {}", - shuffleFallbackCounts.values().stream().mapToLong(v -> v).sum(), - appId); + handleAppHeartbeat(request.getAppHeartbeatRequest()); + break; + + case BatchHeartbeat: + List workerHeartbeats = + request.getBatchHeartbeatRequest().getWorkerHeartbeatsList(); + List appHeartbeats = + request.getBatchHeartbeatRequest().getAppHeartbeatsList(); + LOG.debug( + "Handle batch heartbeat with {} worker heartbeats and {} app heartbeats.", + workerHeartbeats.size(), + appHeartbeats.size()); + for (PbMetaWorkerHeartbeatRequest workerHeartbeat : workerHeartbeats) { + handleWorkerHeartbeat(workerHeartbeat); + } + for (PbMetaAppHeartbeatRequest appHeartbeat : appHeartbeats) { + handleAppHeartbeat(appHeartbeat); } - Map applicationFallbackCounts = - request.getAppHeartbeatRequest().getApplicationFallbackCountsMap(); - metaSystem.updateAppHeartbeatMeta( - appId, - time, - totalWritten, - fileCount, - shuffleCount, - applicationCount, - shuffleFallbackCounts, - applicationFallbackCounts); break; case AppLost: @@ -213,39 +205,7 @@ public org.apache.celeborn.common.protocol.PbMetaRequestResponse handleWriteRequ break; case WorkerHeartbeat: - host = request.getWorkerHeartbeatRequest().getHost(); - rpcPort = request.getWorkerHeartbeatRequest().getRpcPort(); - pushPort = request.getWorkerHeartbeatRequest().getPushPort(); - fetchPort = request.getWorkerHeartbeatRequest().getFetchPort(); - Map pbDiskInfoMap = request.getWorkerHeartbeatRequest().getDisksMap(); - diskInfos = MetaUtil.fromPbDiskInfoMap(pbDiskInfoMap); - 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, - rpcPort, - pushPort, - fetchPort, - replicatePort, - diskInfos); - metaSystem.updateWorkerHeartbeatMeta( - host, - rpcPort, - pushPort, - fetchPort, - replicatePort, - diskInfos, - request.getWorkerHeartbeatRequest().getTime(), - workerStatus, - highWorkload); + handleWorkerHeartbeat(request.getWorkerHeartbeatRequest()); break; case RegisterWorker: @@ -335,6 +295,64 @@ public org.apache.celeborn.common.protocol.PbMetaRequestResponse handleWriteRequ return responseBuilder.build(); } + private void handleWorkerHeartbeat(PbMetaWorkerHeartbeatRequest workerHeartbeat) { + String host = workerHeartbeat.getHost(); + int rpcPort = workerHeartbeat.getRpcPort(); + int pushPort = workerHeartbeat.getPushPort(); + int fetchPort = workerHeartbeat.getFetchPort(); + int replicatePort = workerHeartbeat.getReplicatePort(); + Map diskInfos = MetaUtil.fromPbDiskInfoMap(workerHeartbeat.getDisksMap()); + boolean highWorkload = workerHeartbeat.getHighWorkload(); + WorkerStatus workerStatus; + if (workerHeartbeat.hasWorkerStatus()) { + workerStatus = MetaUtil.fromPbWorkerStatus(workerHeartbeat.getWorkerStatus()); + } else { + workerStatus = WorkerStatus.normalWorkerStatus(); + } + + LOG.debug( + "Handle worker heartbeat for {} {} {} {} {} {}", + host, + rpcPort, + pushPort, + fetchPort, + replicatePort, + diskInfos); + metaSystem.updateWorkerHeartbeatMeta( + host, + rpcPort, + pushPort, + fetchPort, + replicatePort, + diskInfos, + workerHeartbeat.getTime(), + workerStatus, + highWorkload); + } + + private void handleAppHeartbeat(PbMetaAppHeartbeatRequest appHeartbeat) { + String appId = appHeartbeat.getAppId(); + long shuffleCount = appHeartbeat.getShuffleCount(); + LOG.debug("Handle app heartbeat for {} with shuffle count {}", appId, shuffleCount); + Map shuffleFallbackCounts = appHeartbeat.getShuffleFallbackCountsMap(); + if (CollectionUtils.isNotEmpty(shuffleFallbackCounts)) { + LOG.warn( + "{} shuffle fallbacks in app {}", + shuffleFallbackCounts.values().stream().mapToLong(v -> v).sum(), + appId); + } + Map applicationFallbackCounts = appHeartbeat.getApplicationFallbackCountsMap(); + metaSystem.updateAppHeartbeatMeta( + appId, + appHeartbeat.getTime(), + appHeartbeat.getTotalWritten(), + appHeartbeat.getFileCount(), + shuffleCount, + appHeartbeat.getApplicationCount(), + shuffleFallbackCounts, + applicationFallbackCounts); + } + public void writeToSnapShot(File file) throws IOException { try { metaSystem.writeMetaInfoToFile(file); diff --git a/master/src/main/proto/Resource.proto b/master/src/main/proto/Resource.proto index d8fceee1564..7b326003f2a 100644 --- a/master/src/main/proto/Resource.proto +++ b/master/src/main/proto/Resource.proto @@ -42,6 +42,8 @@ enum Type { ReviseLostShuffles = 29; RegisterApplicationInfo = 30; + + BatchHeartbeat = 31; } enum WorkerEventType { @@ -74,6 +76,7 @@ message ResourceRequest { optional ReportWorkerDecommissionRequest reportWorkerDecommissionRequest = 24; optional BatchUnregisterShuffleRequest batchUnregisterShuffleRequest = 25; optional RegisterApplicationInfoRequest registerApplicationInfoRequest = 26; + optional BatchHeartbeatRequest batchHeartbeatRequest = 27; optional ReviseLostShufflesRequest reviseLostShufflesRequest = 102; } @@ -106,6 +109,11 @@ message BatchUnregisterShuffleRequest { repeated string shuffleKeys = 1; } +message BatchHeartbeatRequest { + repeated WorkerHeartbeatRequest workerHeartbeats = 1; + repeated AppHeartbeatRequest appHeartbeats = 2; +} + message HeartbeatInfo { string appId = 1; int64 totalWritten = 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 75ceb951a9d..82abc5834ac 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 @@ -1620,6 +1620,15 @@ private[celeborn] class Master( val transferLeadership = conf.haMasterGracefulShutdownEnabled statusSystem match { case ha: HAMasterMetaManager => + val heartbeatAggregator = ha.getHeartbeatAggregator + if (heartbeatAggregator != null) { + try { + heartbeatAggregator.stop() + } catch { + case e: Exception => + logError("Failed to stop heartbeat aggregator during Master shutdown.", e) + } + } val ratisServer = ha.getRatisServer if (ratisServer != null) { try { diff --git a/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/HeartbeatAggregatorSuiteJ.java b/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/HeartbeatAggregatorSuiteJ.java new file mode 100644 index 00000000000..5e200b774f4 --- /dev/null +++ b/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/HeartbeatAggregatorSuiteJ.java @@ -0,0 +1,230 @@ +/* + * 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.master.clustermeta.ha; + +import java.io.File; +import java.io.IOException; +import java.util.Collections; +import java.util.HashMap; +import java.util.UUID; + +import org.junit.After; +import org.junit.Assert; +import org.junit.Before; +import org.junit.Test; + +import org.apache.celeborn.common.CelebornConf; +import org.apache.celeborn.common.meta.WorkerInfo; +import org.apache.celeborn.common.meta.WorkerStatus; + +public class HeartbeatAggregatorSuiteJ { + private HARaftServer ratisServer; + private HAMasterMetaManager metaSystem; + + @Before + public void init() throws Exception { + CelebornConf conf = new CelebornConf(); + conf.set(CelebornConf.MASTER_HA_HEARTBEAT_BATCH_ENABLED().key(), "true"); + // Use a 1s window (not shorter): on oversubscribed CI runners the handler calls below can + // be scheduled far apart, and a short window would split the offers into several batches. + conf.set(CelebornConf.MASTER_HA_HEARTBEAT_BATCH_INTERVAL().key(), "1s"); + metaSystem = new HAMasterMetaManager(null, conf); + MetaHandler handler = new MetaHandler(metaSystem); + File tmpDir = File.createTempFile("celeborn-ratis-tmp", "for-test-only"); + tmpDir.delete(); + tmpDir.mkdirs(); + conf.set(CelebornConf.HA_MASTER_RATIS_STORAGE_DIR().key(), tmpDir.getAbsolutePath()); + String id = UUID.randomUUID().toString(); + MasterNode masterNode = + new MasterNode.Builder().setNodeId(id).setHost("localhost").setRatisPort(9998).build(); + ratisServer = + HARaftServer.newMasterRatisServer(handler, conf, masterNode, Collections.emptyList()); + metaSystem.setRatisServer(ratisServer); + ratisServer.start(); + waitForLeader(); + } + + private void waitForLeader() throws InterruptedException { + // Wait for isLeaderReady(), not just isLeader(): a newly elected leader rejects writes + // with LeaderNotReadyException until its no-op entry for the new term is committed. + for (int i = 0; i < 100; i++) { + if (ratisServer.isLeader()) { + try { + if (ratisServer + .getServer() + .getDivision(ratisServer.getGroupId()) + .getInfo() + .isLeaderReady()) { + return; + } + } catch (IOException e) { + // Division not available yet; keep polling. + } + } + Thread.sleep(200); + } + Assert.fail("Raft server did not become ready leader in time."); + } + + @After + public void shutdown() { + if (metaSystem.getHeartbeatAggregator() != null) { + metaSystem.getHeartbeatAggregator().stop(); + } + if (ratisServer != null) { + ratisServer.stop(); + } + } + + @Test + public void testBatchHeartbeatAggregation() throws Exception { + Assert.assertNotNull(metaSystem.getHeartbeatAggregator()); + + metaSystem.handleRegisterWorker( + "host1", + 1, + 2, + 3, + 4, + 5, + "networkLocation1", + new HashMap<>(), + new HashMap<>(), + UUID.randomUUID().toString() + "#1"); + Thread.sleep(2000); + Assert.assertEquals(1, metaSystem.workersMap.size()); + long appliedIndexAfterRegister = lastAppliedIndex(); + + // 4 heartbeat offers: 3 for the same worker (collapsed to the newest), 1 for an app. + long time1 = System.currentTimeMillis(); + long time2 = time1 + 10; + metaSystem.handleWorkerHeartbeat( + "host1", + 1, + 2, + 3, + 4, + new HashMap<>(), + new HashMap<>(), + time1, + false, + WorkerStatus.normalWorkerStatus(), + UUID.randomUUID().toString() + "#2"); + metaSystem.handleWorkerHeartbeat( + "host1", + 1, + 2, + 3, + 4, + new HashMap<>(), + new HashMap<>(), + time1, + false, + WorkerStatus.normalWorkerStatus(), + UUID.randomUUID().toString() + "#3"); + metaSystem.handleAppHeartbeat( + "app-1", + 100, + 10, + 1, + 1, + new HashMap<>(), + new HashMap<>(), + time2, + UUID.randomUUID().toString() + "#4"); + metaSystem.handleWorkerHeartbeat( + "host1", + 1, + 2, + 3, + 4, + new HashMap<>(), + new HashMap<>(), + time2, + false, + WorkerStatus.normalWorkerStatus(), + UUID.randomUUID().toString() + "#5"); + + // Wait comfortably past one 1s flush window plus the blocking submit/apply. + Thread.sleep(4000); + + // The 4 offers must have been flushed as (far) fewer raft log entries. + long newEntries = lastAppliedIndex() - appliedIndexAfterRegister; + Assert.assertTrue( + "Expected heartbeats to be merged into few raft log entries, but got " + newEntries, + newEntries >= 1 && newEntries <= 2); + + // The newest heartbeat wins. + WorkerInfo workerInfo = metaSystem.workersMap.values().iterator().next(); + Assert.assertEquals("host1", workerInfo.host()); + Assert.assertEquals(time2, workerInfo.lastHeartbeat()); + Assert.assertEquals(Long.valueOf(time2), metaSystem.appHeartbeatTime.get("app-1")); + } + + @Test + public void testEmptyWindowProducesNoRaftLog() throws Exception { + Thread.sleep(2000); + long appliedIndex = lastAppliedIndex(); + Thread.sleep(500); + Assert.assertEquals(appliedIndex, lastAppliedIndex()); + } + + @Test + public void testDuplicateAppHeartbeatKeepsNewest() throws Exception { + // Duplicate heartbeats of the same app within one window are collapsed to the newest one. + java.util.Map fallback1 = new HashMap<>(); + fallback1.put("shuffle-1", 1L); + java.util.Map fallback2 = new HashMap<>(); + fallback2.put("shuffle-1", 2L); + long time1 = System.currentTimeMillis(); + long time2 = time1 + 10; + metaSystem.handleAppHeartbeat( + "app-merge", + 100, + 10, + 1, + 1, + fallback1, + new HashMap<>(), + time1, + UUID.randomUUID().toString() + "#1"); + metaSystem.handleAppHeartbeat( + "app-merge", + 200, + 20, + 2, + 3, + fallback2, + new HashMap<>(), + time2, + UUID.randomUUID().toString() + "#2"); + + Thread.sleep(4000); + + Assert.assertEquals(200, metaSystem.partitionTotalWritten.sum()); + Assert.assertEquals(20, metaSystem.partitionTotalFileCount.sum()); + Assert.assertEquals(2, metaSystem.shuffleTotalCount.sum()); + Assert.assertEquals(3, metaSystem.applicationTotalCount.sum()); + Assert.assertEquals(Long.valueOf(2), metaSystem.shuffleFallbackCounts.get("shuffle-1")); + Assert.assertEquals(Long.valueOf(time2), metaSystem.appHeartbeatTime.get("app-merge")); + } + + private long lastAppliedIndex() { + return ratisServer.getMasterStateMachine().getLastAppliedTermIndex().getIndex(); + } +} diff --git a/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/MasterStateMachineSuiteJ.java b/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/MasterStateMachineSuiteJ.java index 46a9c68e61b..782d141b798 100644 --- a/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/MasterStateMachineSuiteJ.java +++ b/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/MasterStateMachineSuiteJ.java @@ -552,4 +552,114 @@ public void testInstallSnapshot() .getIndex()); stopRaftServers(raftServers); } + + @Test + public void testBatchHeartbeat() throws InvalidProtocolBufferException { + StateMachine stateMachine = ratisServer.getMasterStateMachine(); + + // Register two workers first, heartbeat apply is a no-op for unregistered workers. + for (int i = 1; i <= 2; i++) { + PbMetaRequest registerRequest = + PbMetaRequest.newBuilder() + .setMetaRequestType( + org.apache.celeborn.common.protocol.PbMetaRequestType.RegisterWorker) + .setRequestId(UUID.randomUUID().toString()) + .setRegisterWorkerRequest( + org.apache.celeborn.common.protocol.PbMetaRegisterWorkerRequest.newBuilder() + .setHost("host" + i) + .setRpcPort(1) + .setPushPort(2) + .setFetchPort(3) + .setReplicatePort(4) + .setInternalPort(5) + .build()) + .build(); + Assert.assertTrue(stateMachine.runCommand(registerRequest, -1).getSuccess()); + } + Assert.assertEquals(2, metaSystem.workersMap.size()); + + long heartbeatTime = System.currentTimeMillis(); + org.apache.celeborn.common.protocol.PbMetaWorkerHeartbeatRequest workerHeartbeat1 = + org.apache.celeborn.common.protocol.PbMetaWorkerHeartbeatRequest.newBuilder() + .setHost("host1") + .setRpcPort(1) + .setPushPort(2) + .setFetchPort(3) + .setReplicatePort(4) + .setTime(heartbeatTime) + .build(); + org.apache.celeborn.common.protocol.PbMetaWorkerHeartbeatRequest workerHeartbeat2 = + org.apache.celeborn.common.protocol.PbMetaWorkerHeartbeatRequest.newBuilder() + .setHost("host2") + .setRpcPort(1) + .setPushPort(2) + .setFetchPort(3) + .setReplicatePort(4) + .setTime(heartbeatTime) + .build(); + org.apache.celeborn.common.protocol.PbMetaAppHeartbeatRequest appHeartbeat = + org.apache.celeborn.common.protocol.PbMetaAppHeartbeatRequest.newBuilder() + .setAppId("app-1") + .setTime(heartbeatTime) + .setTotalWritten(100) + .setFileCount(10) + .build(); + + PbMetaRequest batchRequest = + PbMetaRequest.newBuilder() + .setMetaRequestType( + org.apache.celeborn.common.protocol.PbMetaRequestType.BatchHeartbeat) + .setRequestId(UUID.randomUUID().toString()) + .setBatchHeartbeatRequest( + org.apache.celeborn.common.protocol.PbBatchHeartbeatRequest.newBuilder() + .addWorkerHeartbeats(workerHeartbeat1) + .addWorkerHeartbeats(workerHeartbeat2) + .addAppHeartbeats(appHeartbeat) + .build()) + .build(); + + PbMetaRequestResponse response = stateMachine.runCommand(batchRequest, -1); + Assert.assertTrue(response.getSuccess()); + + // Each worker heartbeat must be applied with its own time field (see CELEBORN-2399). + for (WorkerInfo workerInfo : metaSystem.workersMap.values()) { + Assert.assertEquals(heartbeatTime, workerInfo.lastHeartbeat()); + } + Assert.assertEquals(Long.valueOf(heartbeatTime), metaSystem.appHeartbeatTime.get("app-1")); + + // The batched log entry written by the leader (ResourceRequest) must stay wire-compatible + // with the apply side (PbMetaRequest). + ResourceProtos.ResourceRequest resourceRequest = + ResourceProtos.ResourceRequest.newBuilder() + .setCmdType(ResourceProtos.Type.BatchHeartbeat) + .setRequestId(UUID.randomUUID().toString()) + .setBatchHeartbeatRequest( + ResourceProtos.BatchHeartbeatRequest.newBuilder() + .addWorkerHeartbeats( + ResourceProtos.WorkerHeartbeatRequest.newBuilder() + .setHost("host1") + .setRpcPort(1) + .setPushPort(2) + .setFetchPort(3) + .setReplicatePort(4) + .setTime(heartbeatTime) + .build()) + .addAppHeartbeats( + ResourceProtos.AppHeartbeatRequest.newBuilder() + .setAppId("app-1") + .setTime(heartbeatTime) + .setTotalWritten(100) + .setFileCount(10) + .build()) + .build()) + .build(); + PbMetaRequest parsed = PbMetaRequest.parseFrom(resourceRequest.toByteString()); + Assert.assertEquals( + org.apache.celeborn.common.protocol.PbMetaRequestType.BatchHeartbeat, + parsed.getMetaRequestType()); + Assert.assertEquals(1, parsed.getBatchHeartbeatRequest().getWorkerHeartbeatsCount()); + Assert.assertEquals(1, parsed.getBatchHeartbeatRequest().getAppHeartbeatsCount()); + Assert.assertEquals( + heartbeatTime, parsed.getBatchHeartbeatRequest().getWorkerHeartbeats(0).getTime()); + } } diff --git a/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/RatisBaseSuiteJ.java b/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/RatisBaseSuiteJ.java index b727758e4af..3a03c636189 100644 --- a/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/RatisBaseSuiteJ.java +++ b/master/src/test/java/org/apache/celeborn/service/deploy/master/clustermeta/ha/RatisBaseSuiteJ.java @@ -28,11 +28,12 @@ public class RatisBaseSuiteJ { HARaftServer ratisServer; + HAMasterMetaManager metaSystem; @Before public void init() throws Exception { CelebornConf conf = new CelebornConf(); - HAMasterMetaManager metaSystem = new HAMasterMetaManager(null, conf); + metaSystem = new HAMasterMetaManager(null, conf); MetaHandler handler = new MetaHandler(metaSystem); File tmpDir1 = File.createTempFile("celeborn-ratis-tmp", "for-test-only"); tmpDir1.delete();