This is an automated email from the ASF dual-hosted git repository. tanxinyu pushed a commit to branch jira3512 in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit d979d9d8c0a7c5935ded6746c5a8ddc794b2a2cd Author: LebronAl <[email protected]> AuthorDate: Fri Jun 17 13:45:32 2022 +0800 finish --- .../org/apache/iotdb/consensus/IConsensus.java | 2 ++ .../multileader/MultiLeaderConsensus.java | 5 +++++ .../iotdb/consensus/ratis/RatisConsensus.java | 12 +++++++++++ .../consensus/standalone/StandAloneConsensus.java | 6 ++++++ .../iotdb/consensus/ratis/RatisConsensusTest.java | 4 ---- .../service/thrift/impl/InternalServiceImpl.java | 25 +++++++++++++++++++++- thrift-commons/src/main/thrift/common.thrift | 5 +++-- 7 files changed, 52 insertions(+), 7 deletions(-) diff --git a/consensus/src/main/java/org/apache/iotdb/consensus/IConsensus.java b/consensus/src/main/java/org/apache/iotdb/consensus/IConsensus.java index f7c195eefc..60982cc49b 100644 --- a/consensus/src/main/java/org/apache/iotdb/consensus/IConsensus.java +++ b/consensus/src/main/java/org/apache/iotdb/consensus/IConsensus.java @@ -63,4 +63,6 @@ public interface IConsensus { boolean isLeader(ConsensusGroupId groupId); Peer getLeader(ConsensusGroupId groupId); + + List<ConsensusGroupId> getAllConsensusGroupIds(); } diff --git a/consensus/src/main/java/org/apache/iotdb/consensus/multileader/MultiLeaderConsensus.java b/consensus/src/main/java/org/apache/iotdb/consensus/multileader/MultiLeaderConsensus.java index fbd9b280f3..b9e762a1db 100644 --- a/consensus/src/main/java/org/apache/iotdb/consensus/multileader/MultiLeaderConsensus.java +++ b/consensus/src/main/java/org/apache/iotdb/consensus/multileader/MultiLeaderConsensus.java @@ -236,6 +236,11 @@ public class MultiLeaderConsensus implements IConsensus { return new Peer(groupId, thisNode); } + @Override + public List<ConsensusGroupId> getAllConsensusGroupIds() { + return new ArrayList<>(stateMachineMap.keySet()); + } + public MultiLeaderServerImpl getImpl(ConsensusGroupId groupId) { return stateMachineMap.get(groupId); } diff --git a/consensus/src/main/java/org/apache/iotdb/consensus/ratis/RatisConsensus.java b/consensus/src/main/java/org/apache/iotdb/consensus/ratis/RatisConsensus.java index 8687dba132..33a523eb06 100644 --- a/consensus/src/main/java/org/apache/iotdb/consensus/ratis/RatisConsensus.java +++ b/consensus/src/main/java/org/apache/iotdb/consensus/ratis/RatisConsensus.java @@ -530,6 +530,18 @@ class RatisConsensus implements IConsensus { return new Peer(groupId, leaderEndpoint); } + @Override + public List<ConsensusGroupId> getAllConsensusGroupIds() { + List<ConsensusGroupId> ids = new ArrayList<>(); + server + .getGroupIds() + .forEach( + groupId -> { + ids.add(Utils.fromRaftGroupIdToConsensusGroupId(groupId)); + }); + return ids; + } + @Override public ConsensusGenericResponse triggerSnapshot(ConsensusGroupId groupId) { RaftGroupId raftGroupId = Utils.fromConsensusGroupIdToRaftGroupId(groupId); diff --git a/consensus/src/main/java/org/apache/iotdb/consensus/standalone/StandAloneConsensus.java b/consensus/src/main/java/org/apache/iotdb/consensus/standalone/StandAloneConsensus.java index 1aa568416b..d87e27837f 100644 --- a/consensus/src/main/java/org/apache/iotdb/consensus/standalone/StandAloneConsensus.java +++ b/consensus/src/main/java/org/apache/iotdb/consensus/standalone/StandAloneConsensus.java @@ -44,6 +44,7 @@ import java.io.IOException; import java.nio.file.DirectoryStream; import java.nio.file.Files; import java.nio.file.Path; +import java.util.ArrayList; import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; @@ -218,6 +219,11 @@ class StandAloneConsensus implements IConsensus { return new Peer(groupId, thisNode); } + @Override + public List<ConsensusGroupId> getAllConsensusGroupIds() { + return new ArrayList<>(stateMachineMap.keySet()); + } + private String buildPeerDir(ConsensusGroupId groupId) { return storageDir + File.separator + groupId.getType().getValue() + "_" + groupId.getId(); } diff --git a/consensus/src/test/java/org/apache/iotdb/consensus/ratis/RatisConsensusTest.java b/consensus/src/test/java/org/apache/iotdb/consensus/ratis/RatisConsensusTest.java index 456a27fdde..b44d96460d 100644 --- a/consensus/src/test/java/org/apache/iotdb/consensus/ratis/RatisConsensusTest.java +++ b/consensus/src/test/java/org/apache/iotdb/consensus/ratis/RatisConsensusTest.java @@ -119,10 +119,6 @@ public class RatisConsensusTest { } } - private int getLeaderOrdinal() { - return servers.get(0).getLeader(gid).getEndpoint().port - 6000; - } - @Test public void basicConsensus3Copy() throws Exception { servers.get(0).addConsensusGroup(group.getGroupId(), group.getPeers()); diff --git a/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/InternalServiceImpl.java b/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/InternalServiceImpl.java index 17efadbf68..557f3f3fa9 100644 --- a/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/InternalServiceImpl.java +++ b/server/src/main/java/org/apache/iotdb/db/service/thrift/impl/InternalServiceImpl.java @@ -95,7 +95,9 @@ import org.slf4j.LoggerFactory; import java.util.ArrayList; import java.util.Arrays; +import java.util.HashMap; import java.util.List; +import java.util.Map; import java.util.Random; import java.util.stream.Collectors; @@ -305,7 +307,7 @@ public class InternalServiceImpl implements InternalService.Iface { @Override public THeartbeatResp getHeartBeat(THeartbeatReq req) throws TException { - THeartbeatResp resp = new THeartbeatResp(req.getHeartbeatTimestamp()); + THeartbeatResp resp = new THeartbeatResp(req.getHeartbeatTimestamp(), getJudgedLeaders()); Random whetherToGetMetric = new Random(); if (MetricConfigDescriptor.getInstance().getMetricConfig().getEnableMetric() && whetherToGetMetric.nextDouble() < loadBalanceThreshold) { @@ -327,6 +329,27 @@ public class InternalServiceImpl implements InternalService.Iface { return resp; } + private Map<TConsensusGroupId, Boolean> getJudgedLeaders() { + Map<TConsensusGroupId, Boolean> result = new HashMap<>(); + DataRegionConsensusImpl.getInstance() + .getAllConsensusGroupIds() + .forEach( + groupId -> { + result.put( + groupId.convertToTConsensusGroupId(), + DataRegionConsensusImpl.getInstance().isLeader(groupId)); + }); + SchemaRegionConsensusImpl.getInstance() + .getAllConsensusGroupIds() + .forEach( + groupId -> { + result.put( + groupId.convertToTConsensusGroupId(), + SchemaRegionConsensusImpl.getInstance().isLeader(groupId)); + }); + return result; + } + private long getMemory(String gaugeName) { long result = 0; try { diff --git a/thrift-commons/src/main/thrift/common.thrift b/thrift-commons/src/main/thrift/common.thrift index e314a6327a..c381d5fb2d 100644 --- a/thrift-commons/src/main/thrift/common.thrift +++ b/thrift-commons/src/main/thrift/common.thrift @@ -83,8 +83,9 @@ struct THeartbeatReq { struct THeartbeatResp { 1: required i64 heartbeatTimestamp - 2: optional i16 cpu - 3: optional i16 memory + 2: required map<TConsensusGroupId, bool> judgedLeaders + 3: optional i16 cpu + 4: optional i16 memory } struct TDataNodeInfo {
