Fold ConfigNode priority reconcile into EventService; drop iterationIndex threading - HeartbeatService bumps heartbeatCounter once at the tail of heartbeatLoopBody; the gen* methods read heartbeatCounter.get() directly again instead of taking an iterationIndex parameter. - Remove the standalone ConfigRegionPriorityBalancer. Priority reconciliation now lives in EventService, fired from checkAndBroadcastNodeStatisticsChangeEventIfNecessary exactly when ConfigNode statistics change. It is leader-gated, filtered to ConfigNode peers, and only pushes peers whose priority bucket actually moved. EventService now takes the IManager to reach the consensus impl.
diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/LoadManager.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/LoadManager.java index 1ef9100..2706226 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/LoadManager.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/LoadManager.java
@@ -34,7 +34,6 @@ import org.apache.iotdb.confignode.exception.NoAvailableRegionGroupException; import org.apache.iotdb.confignode.exception.NotEnoughDataNodeException; import org.apache.iotdb.confignode.manager.IManager; -import org.apache.iotdb.confignode.manager.load.balancer.ConfigRegionPriorityBalancer; import org.apache.iotdb.confignode.manager.load.balancer.PartitionBalancer; import org.apache.iotdb.confignode.manager.load.balancer.RegionBalancer; import org.apache.iotdb.confignode.manager.load.balancer.RouteBalancer; @@ -87,11 +86,10 @@ setHeartbeatService(configManager, loadCache); this.statisticsService = new StatisticsService(loadCache); this.topologyService = new TopologyService(configManager, loadCache::updateTopology); - this.eventService = new EventService(loadCache); + this.eventService = new EventService(configManager, loadCache); this.eventService.register(configManager.getPipeManager().getPipeRuntimeCoordinator()); this.eventService.register(routeBalancer); this.eventService.register(topologyService); - this.eventService.register(new ConfigRegionPriorityBalancer(configManager)); } protected void setHeartbeatService(IManager configManager, LoadCache loadCache) {
diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/balancer/ConfigRegionPriorityBalancer.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/balancer/ConfigRegionPriorityBalancer.java deleted file mode 100644 index dc9d1d5..0000000 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/balancer/ConfigRegionPriorityBalancer.java +++ /dev/null
@@ -1,114 +0,0 @@ -/* - * 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.iotdb.confignode.manager.load.balancer; - -import org.apache.iotdb.common.rpc.thrift.TConfigNodeLocation; -import org.apache.iotdb.commons.cluster.NodeStatus; -import org.apache.iotdb.confignode.i18n.ManagerMessages; -import org.apache.iotdb.confignode.manager.IManager; -import org.apache.iotdb.confignode.manager.consensus.ConsensusManager; -import org.apache.iotdb.confignode.manager.load.cache.node.NodeStatistics; -import org.apache.iotdb.confignode.manager.load.subscriber.IClusterStatusSubscriber; -import org.apache.iotdb.confignode.manager.load.subscriber.NodeStatisticsChangeEvent; -import org.apache.iotdb.consensus.exception.ConsensusException; - -import org.apache.tsfile.utils.Pair; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import java.util.HashMap; -import java.util.Map; -import java.util.OptionalInt; -import java.util.Set; -import java.util.stream.Collectors; - -/** - * Reacts to {@link NodeStatisticsChangeEvent} on the ConfigNode leader and pushes Ratis peer - * priorities for the ConfigRegion so candidacy in leader elections reflects the latest {@link - * NodeStatus} of every ConfigNode (see {@link NodeStatus#priorityForStatus}). - * - * <p>Only fires when this ConfigNode is the ConfigRegion leader, since {@code setConfiguration} - * must be applied through the leader. Filters events to ConfigNode peers (DataNode/AINode entries - * are skipped) and accumulates only transitions whose priority bucket actually moves before issuing - * a single batched reconfiguration call. - */ -public class ConfigRegionPriorityBalancer implements IClusterStatusSubscriber { - - private static final Logger LOGGER = LoggerFactory.getLogger(ConfigRegionPriorityBalancer.class); - - private final IManager configManager; - - public ConfigRegionPriorityBalancer(IManager configManager) { - this.configManager = configManager; - } - - @Override - public void onNodeStatisticsChanged(NodeStatisticsChangeEvent event) { - if (!configManager.getConsensusManager().isLeader()) { - return; - } - Set<Integer> configNodeIds = - configManager.getNodeManager().getRegisteredConfigNodes().stream() - .map(TConfigNodeLocation::getConfigNodeId) - .collect(Collectors.toSet()); - - Map<Integer, Integer> desired = new HashMap<>(); - for (Map.Entry<Integer, Pair<NodeStatistics, NodeStatistics>> entry : - event.getDifferentNodeStatisticsMap().entrySet()) { - int nodeId = entry.getKey(); - if (!configNodeIds.contains(nodeId)) { - continue; - } - NodeStatistics previous = entry.getValue().getLeft(); - NodeStatistics current = entry.getValue().getRight(); - if (current == null) { - // Node disappeared from the cache; peer removal flows through addConfigNodePeer / - // removeConfigNodePeer instead — no priority push to issue here. - continue; - } - OptionalInt newPriority = - NodeStatus.priorityForStatus(current.getStatus(), current.getStatusReason()); - if (!newPriority.isPresent()) { - // Transient (Unknown / Removing / manual ReadOnly) — leave the priority as-is. - continue; - } - OptionalInt oldPriority = - previous == null - ? OptionalInt.empty() - : NodeStatus.priorityForStatus(previous.getStatus(), previous.getStatusReason()); - if (oldPriority.isPresent() && oldPriority.getAsInt() == newPriority.getAsInt()) { - // The status moved but the priority bucket did not — no point churning the group config. - continue; - } - desired.put(nodeId, newPriority.getAsInt()); - } - if (desired.isEmpty()) { - return; - } - try { - configManager - .getConsensusManager() - .getConsensusImpl() - .reconfigurePeerPriorities(ConsensusManager.DEFAULT_CONSENSUS_GROUP_ID, desired); - } catch (ConsensusException e) { - LOGGER.warn(ManagerMessages.RECONFIGURE_PEER_PRIORITIES_FAILED, desired, e); - } - } -}
diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/service/EventService.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/service/EventService.java index 762e6b1..9828d90 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/service/EventService.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/service/EventService.java
@@ -19,13 +19,17 @@ package org.apache.iotdb.confignode.manager.load.service; +import org.apache.iotdb.common.rpc.thrift.TConfigNodeLocation; import org.apache.iotdb.common.rpc.thrift.TConsensusGroupId; +import org.apache.iotdb.commons.cluster.NodeStatus; import org.apache.iotdb.commons.concurrent.IoTDBThreadPoolFactory; import org.apache.iotdb.commons.concurrent.ThreadName; import org.apache.iotdb.commons.concurrent.threadpool.ScheduledExecutorUtil; import org.apache.iotdb.commons.utils.TestOnly; import org.apache.iotdb.confignode.conf.ConfigNodeDescriptor; import org.apache.iotdb.confignode.i18n.ManagerMessages; +import org.apache.iotdb.confignode.manager.IManager; +import org.apache.iotdb.confignode.manager.consensus.ConsensusManager; import org.apache.iotdb.confignode.manager.load.cache.LoadCache; import org.apache.iotdb.confignode.manager.load.cache.consensus.ConsensusGroupStatistics; import org.apache.iotdb.confignode.manager.load.cache.node.NodeStatistics; @@ -34,6 +38,7 @@ import org.apache.iotdb.confignode.manager.load.subscriber.IClusterStatusSubscriber; import org.apache.iotdb.confignode.manager.load.subscriber.NodeStatisticsChangeEvent; import org.apache.iotdb.confignode.manager.load.subscriber.RegionGroupStatisticsChangeEvent; +import org.apache.iotdb.consensus.exception.ConsensusException; import com.google.common.eventbus.AsyncEventBus; import com.google.common.eventbus.EventBus; @@ -42,13 +47,17 @@ import org.slf4j.LoggerFactory; import java.util.Collections; +import java.util.HashMap; import java.util.List; import java.util.Map; import java.util.Objects; +import java.util.OptionalInt; +import java.util.Set; import java.util.TreeMap; import java.util.concurrent.Future; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; /** * EventService periodically check statistics and broadcast corresponding change event if necessary. @@ -68,6 +77,7 @@ IoTDBThreadPoolFactory.newSingleThreadScheduledExecutor( ThreadName.CONFIG_NODE_EVENT_SERVICE.getName()); + private final IManager configManager; private final LoadCache loadCache; private final Map<Integer, NodeStatistics> previousNodeStatisticsMap; private final Map<TConsensusGroupId, RegionGroupStatistics> previousRegionGroupStatisticsMap; @@ -75,7 +85,8 @@ previousConsensusGroupStatisticsMap; private final EventBus eventPublisher; - public EventService(LoadCache loadCache) { + public EventService(IManager configManager, LoadCache loadCache) { + this.configManager = configManager; this.loadCache = loadCache; this.previousNodeStatisticsMap = new TreeMap<>(); this.previousRegionGroupStatisticsMap = new TreeMap<>(); @@ -154,6 +165,55 @@ if (!differentNodeStatisticsMap.isEmpty()) { eventPublisher.post(new NodeStatisticsChangeEvent(differentNodeStatisticsMap)); recordNodeStatistics(differentNodeStatisticsMap); + reconcileConfigNodePriorities(differentNodeStatisticsMap); + } + } + + /** + * Demote a ConfigNode's ConfigRegion election priority when its {@link NodeStatus} degrades (see + * {@link NodeStatus#priorityForStatus}). Runs only on the leader and only pushes peers whose + * priority bucket actually moved. + */ + private void reconcileConfigNodePriorities( + Map<Integer, Pair<NodeStatistics, NodeStatistics>> differentNodeStatisticsMap) { + ConsensusManager consensusManager = configManager.getConsensusManager(); + if (consensusManager == null || !consensusManager.isLeader()) { + return; + } + Set<Integer> configNodeIds = + configManager.getNodeManager().getRegisteredConfigNodes().stream() + .map(TConfigNodeLocation::getConfigNodeId) + .collect(Collectors.toSet()); + Map<Integer, Integer> desired = new HashMap<>(); + differentNodeStatisticsMap.forEach( + (nodeId, change) -> { + NodeStatistics current = change.getRight(); + if (!configNodeIds.contains(nodeId) || current == null) { + return; + } + OptionalInt newPriority = + NodeStatus.priorityForStatus(current.getStatus(), current.getStatusReason()); + if (!newPriority.isPresent()) { + return; + } + NodeStatistics previous = change.getLeft(); + OptionalInt oldPriority = + previous == null + ? OptionalInt.empty() + : NodeStatus.priorityForStatus(previous.getStatus(), previous.getStatusReason()); + if (!oldPriority.isPresent() || oldPriority.getAsInt() != newPriority.getAsInt()) { + desired.put(nodeId, newPriority.getAsInt()); + } + }); + if (desired.isEmpty()) { + return; + } + try { + consensusManager + .getConsensusImpl() + .reconfigurePeerPriorities(ConsensusManager.DEFAULT_CONSENSUS_GROUP_ID, desired); + } catch (ConsensusException e) { + LOGGER.warn(ManagerMessages.RECONFIGURE_PEER_PRIORITIES_FAILED, desired, e); } }
diff --git a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/service/HeartbeatService.java b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/service/HeartbeatService.java index 915fbe9..73e28b4 100644 --- a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/service/HeartbeatService.java +++ b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/load/service/HeartbeatService.java
@@ -137,36 +137,26 @@ .ifPresent( consensusManager -> { if (getConsensusManager().isLeader()) { - // The counter is bumped exactly once per loop iteration so every gen* call in - // this body observes the same value, and the leader's own sampling fires on the - // same iterations DataNode/AINode load sampling does. - long iterationIndex = heartbeatCounter.getAndIncrement(); // Send heartbeat requests to all the registered ConfigNodes pingRegisteredConfigNodes( - genConfigNodeHeartbeatReq(iterationIndex), - getNodeManager().getRegisteredConfigNodes()); + genConfigNodeHeartbeatReq(), getNodeManager().getRegisteredConfigNodes()); // Send heartbeat requests to all the registered DataNodes pingRegisteredDataNodes( - genHeartbeatReq(iterationIndex), - getNodeManager().getRegisteredDataNodes(), - iterationIndex); + genHeartbeatReq(), getNodeManager().getRegisteredDataNodes()); // Send heartbeat requests to all the registered AINodes - pingRegisteredAINodes( - genAIHeartbeatReq(iterationIndex), getNodeManager().getRegisteredAINodes()); - // Sample free-space on the same cadence DataNode samples its load. Runs after - // the async heartbeat dispatches so the OS call does not delay fanout. DiskCrash - // is observed passively by the Ratis write-path, not polled here. Cross-peer - // priority reconciliation is event-driven (see ConfigRegionPriorityBalancer). - if (iterationIndex % LOAD_SAMPLING_INTERVAL == 0) { + pingRegisteredAINodes(genAIHeartbeatReq(), getNodeManager().getRegisteredAINodes()); + // Self-check free-space after the async fanout so the OS call doesn't delay it. + if (heartbeatCounter.get() % LOAD_SAMPLING_INTERVAL == 0) { DiskChecker.checkFreeRatioAndApply( ConfigNodeDescriptor.getInstance().getConf().getCriticalDirs(), CommonDescriptor.getInstance().getConfig().getDiskSpaceWarningThreshold()); } + heartbeatCounter.getAndIncrement(); } }); } - protected TDataNodeHeartbeatReq genHeartbeatReq(long iterationIndex) { + protected TDataNodeHeartbeatReq genHeartbeatReq() { /* Generate heartbeat request */ TDataNodeHeartbeatReq heartbeatReq = new TDataNodeHeartbeatReq(); heartbeatReq.setHeartbeatTimestamp(System.nanoTime()); @@ -177,7 +167,7 @@ .getLogicalClock(ConfigNodeInfo.CONFIG_REGION_ID)); // Always sample RegionGroups' leadership as the Region heartbeat heartbeatReq.setNeedJudgeLeader(true); - heartbeatReq.setNeedSamplingLoad(iterationIndex % LOAD_SAMPLING_INTERVAL == 0); + heartbeatReq.setNeedSamplingLoad(heartbeatCounter.get() % LOAD_SAMPLING_INTERVAL == 0); Pair<Long, Long> schemaQuotaRemain = configManager.getClusterSchemaManager().getSchemaQuotaRemain(); heartbeatReq.setTimeSeriesQuotaRemain(schemaQuotaRemain.left); @@ -185,7 +175,7 @@ // We collect pipe meta in every 100 heartbeat loop heartbeatReq.setNeedPipeMetaList( !PipeConfig.getInstance().isSeperatedPipeHeartbeatEnabled() - && iterationIndex + && heartbeatCounter.get() % PipeConfig.getInstance() .getPipeHeartbeatIntervalSecondsForCollectingPipeMeta() == 0); @@ -196,7 +186,7 @@ } // We broadcast region operations list every 100 heartbeat loops - if (iterationIndex % 100 == 0) { + if (heartbeatCounter.get() % 100 == 0) { heartbeatReq.setCurrentRegionOperations( configManager.getProcedureManager().getRegionOperationConsensusIds()); } @@ -204,8 +194,7 @@ return heartbeatReq; } - private void addConfigNodeLocationsToReq( - int dataNodeId, TDataNodeHeartbeatReq req, long iterationIndex) { + private void addConfigNodeLocationsToReq(int dataNodeId, TDataNodeHeartbeatReq req) { Set<TEndPoint> confirmedConfigNodes = loadCache.getConfirmedConfigNodeEndPoints(dataNodeId); Set<TEndPoint> actualConfigNodes = getNodeManager().getRegisteredConfigNodes().stream() @@ -220,23 +209,23 @@ 4. At this point, because actualConfigNodes and confirmedConfigNodes are identical, the ConfigNode list is not re-sent to the DataNode. */ if (!actualConfigNodes.equals(confirmedConfigNodes) - || iterationIndex % configNodeListPeriodicallySyncInterval == 0) { + || heartbeatCounter.get() % configNodeListPeriodicallySyncInterval == 0) { req.setConfigNodeEndPoints(actualConfigNodes); } } - protected TConfigNodeHeartbeatReq genConfigNodeHeartbeatReq(long iterationIndex) { + protected TConfigNodeHeartbeatReq genConfigNodeHeartbeatReq() { TConfigNodeHeartbeatReq req = new TConfigNodeHeartbeatReq(); req.setTimestamp(System.nanoTime()); - req.setNeedSamplingLoad(iterationIndex % LOAD_SAMPLING_INTERVAL == 0); + req.setNeedSamplingLoad(heartbeatCounter.get() % LOAD_SAMPLING_INTERVAL == 0); return req; } - private TAIHeartbeatReq genAIHeartbeatReq(long iterationIndex) { + private TAIHeartbeatReq genAIHeartbeatReq() { /* Generate heartbeat request */ TAIHeartbeatReq heartbeatReq = new TAIHeartbeatReq(); heartbeatReq.setHeartbeatTimestamp(System.nanoTime()); - heartbeatReq.setNeedSamplingLoad(iterationIndex % LOAD_SAMPLING_INTERVAL == 0); + heartbeatReq.setNeedSamplingLoad(heartbeatCounter.get() % LOAD_SAMPLING_INTERVAL == 0); return heartbeatReq; } @@ -273,9 +262,7 @@ * @param registeredDataNodes DataNodes that registered in cluster */ private void pingRegisteredDataNodes( - TDataNodeHeartbeatReq heartbeatReq, - List<TDataNodeConfiguration> registeredDataNodes, - long iterationIndex) { + TDataNodeHeartbeatReq heartbeatReq, List<TDataNodeConfiguration> registeredDataNodes) { // Send heartbeat requests for (TDataNodeConfiguration dataNodeInfo : registeredDataNodes) { int dataNodeId = dataNodeInfo.getLocation().getDataNodeId(); @@ -294,7 +281,7 @@ configManager.getClusterSchemaManager()::updateDeviceUsage, configManager.getPipeManager().getPipeRuntimeCoordinator()); configManager.getClusterQuotaManager().updateSpaceQuotaUsage(); - addConfigNodeLocationsToReq(dataNodeId, heartbeatReq, iterationIndex); + addConfigNodeLocationsToReq(dataNodeId, heartbeatReq); AsyncDataNodeHeartbeatClientPool.getInstance() .getDataNodeHeartBeat( dataNodeInfo.getLocation().getInternalEndPoint(), heartbeatReq, handler);