This is an automated email from the ASF dual-hosted git repository. lta pushed a commit to branch cluster_multi_raft in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit c4a9dcde75b8beba88bf84de2746a5cbf658b47c Author: lta <[email protected]> AuthorDate: Mon Dec 28 17:31:28 2020 +0800 fix a bug of multi-raft --- .../main/java/org/apache/iotdb/cluster/ClusterMain.java | 9 ++++++--- .../org/apache/iotdb/cluster/config/ClusterConfig.java | 2 +- .../org/apache/iotdb/cluster/log/LogDispatcher.java | 1 + .../iotdb/cluster/log/catchup/LogCatchUpTask.java | 2 ++ .../org/apache/iotdb/cluster/metadata/CMManager.java | 4 ++-- .../apache/iotdb/cluster/partition/PartitionGroup.java | 2 +- .../cluster/query/aggregate/ClusterAggregator.java | 1 + .../cluster/query/reader/ClusterReaderFactory.java | 2 ++ .../iotdb/cluster/server/heartbeat/HeartbeatThread.java | 2 ++ .../apache/iotdb/cluster/server/member/RaftMember.java | 1 + .../main/java/org/apache/iotdb/db/conf/IoTDBConfig.java | 2 +- .../java/org/apache/iotdb/rpc/RpcTransportFactory.java | 2 +- thrift/src/main/thrift/cluster.thrift | 17 ++++++++--------- 13 files changed, 29 insertions(+), 18 deletions(-) diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/ClusterMain.java b/cluster/src/main/java/org/apache/iotdb/cluster/ClusterMain.java index 14bedfd..3a67333 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/ClusterMain.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/ClusterMain.java @@ -282,7 +282,9 @@ public class ClusterMain { // nodes evenly, and use default strategy for other groups SlotPartitionTable.setSlotStrategy(new SlotStrategy() { SlotStrategy defaultStrategy = new SlotStrategy.DefaultStrategy(); - int k = 2; + int k = ClusterDescriptor.getInstance().getConfig().getMultiRaftFactor() * ClusterDescriptor + .getInstance().getConfig().getSeedNodeUrls().size(); + @Override public int calculateSlotByTime(String storageGroupName, long timestamp, int maxSlotNum) { int sgSerialNum = extractSerialNumInSGName(storageGroupName) % k; @@ -305,11 +307,12 @@ public class ClusterMain { } private int extractSerialNumInSGName(String storageGroupName) { - String[] s = storageGroupName.split("\\."); +// String[] s = storageGroupName.split("\\."); + String[] s = storageGroupName.split("_"); if (s.length != 2) { return -1; } - s[1] = s[1].substring(4); +// s[1] = s[1].substring(4); try { return Integer.parseInt(s[1]); } catch (NumberFormatException e) { diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterConfig.java b/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterConfig.java index 306b081..cc19a52 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterConfig.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/config/ClusterConfig.java @@ -56,7 +56,7 @@ public class ClusterConfig { private boolean useAsyncApplier = true; - private int connectionTimeoutInMS = 20 * 1000; + private int connectionTimeoutInMS = 20_1000; private int readOperationTimeoutMS = 30_1000; diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/log/LogDispatcher.java b/cluster/src/main/java/org/apache/iotdb/cluster/log/LogDispatcher.java index 21820d5..9cab9fa 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/log/LogDispatcher.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/log/LogDispatcher.java @@ -301,6 +301,7 @@ public class LogDispatcher { private AppendEntriesRequest prepareRequest(List<ByteBuffer> logList, List<SendLogRequest> currBatch, int firstIndex) { AppendEntriesRequest request = new AppendEntriesRequest(); + request.setRaftId(member.getRaftGroupId()); if (member.getHeader() != null) { request.setHeader(member.getHeader()); diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/log/catchup/LogCatchUpTask.java b/cluster/src/main/java/org/apache/iotdb/cluster/log/catchup/LogCatchUpTask.java index 8b69884..317e06f 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/log/catchup/LogCatchUpTask.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/log/catchup/LogCatchUpTask.java @@ -86,6 +86,7 @@ public class LogCatchUpTask implements Callable<Boolean> { void doLogCatchUp() throws TException, InterruptedException, LeaderUnknownException { AppendEntryRequest request = new AppendEntryRequest(); + request.setRaftId(raftId); if (raftMember.getHeader() != null) { request.setHeader(raftMember.getHeader()); } @@ -170,6 +171,7 @@ public class LogCatchUpTask implements Callable<Boolean> { private AppendEntriesRequest prepareRequest(List<ByteBuffer> logList, int startPos) { AppendEntriesRequest request = new AppendEntriesRequest(); + request.setRaftId(raftId); if (raftMember.getHeader() != null) { request.setHeader(raftMember.getHeader()); diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/metadata/CMManager.java b/cluster/src/main/java/org/apache/iotdb/cluster/metadata/CMManager.java index df772a3..01817a4 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/metadata/CMManager.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/metadata/CMManager.java @@ -938,8 +938,8 @@ public class CMManager extends MManager { List<Node> coordinatedNodes = QueryCoordinator.getINSTANCE().reorderNodes(partitionGroup); for (Node node : coordinatedNodes) { try { - List<PartialPath> paths = getMatchedPaths(node, partitionGroup.getHeader(), partitionGroup.getId(), pathsToQuery, - withAlias); + List<PartialPath> paths = getMatchedPaths(node, partitionGroup.getHeader(), + partitionGroup.getId(), pathsToQuery, withAlias); if (logger.isDebugEnabled()) { logger.debug("{}: get matched paths of {} and other {} paths from {} in {}, result {}", metaGroupMember.getName(), pathsToQuery.get(0), pathsToQuery.size() - 1, node, diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/partition/PartitionGroup.java b/cluster/src/main/java/org/apache/iotdb/cluster/partition/PartitionGroup.java index 5ab4275..2a562ac 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/partition/PartitionGroup.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/partition/PartitionGroup.java @@ -66,7 +66,7 @@ public class PartitionGroup extends ArrayList<Node> { @Override public int hashCode() { - return Objects.hash(id, getHeader()); + return Objects.hash(id, super.hashCode()); } public Node getHeader() { diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/query/aggregate/ClusterAggregator.java b/cluster/src/main/java/org/apache/iotdb/cluster/query/aggregate/ClusterAggregator.java index c2b5974..6fb81b5 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/query/aggregate/ClusterAggregator.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/query/aggregate/ClusterAggregator.java @@ -158,6 +158,7 @@ public class ClusterAggregator { QueryContext context, boolean ascending) throws StorageEngineException { GetAggrResultRequest request = new GetAggrResultRequest(); + request.setRaftId(partitionGroup.getId()); request.setPath(path.getFullPath()); request.setAggregations(aggregations); request.setDataTypeOrdinal(dataType.ordinal()); diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/ClusterReaderFactory.java b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/ClusterReaderFactory.java index d99bbbf..c5f6ee6 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/ClusterReaderFactory.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/query/reader/ClusterReaderFactory.java @@ -339,6 +339,7 @@ public class ClusterReaderFactory { Set<String> deviceMeasurements, PartitionGroup partitionGroup, QueryContext context, boolean ascending) { SingleSeriesQueryRequest request = new SingleSeriesQueryRequest(); + request.setRaftId(partitionGroup.getId()); if (timeFilter != null) { request.setTimeFilterBytes(SerializeUtils.serializeFilter(timeFilter)); } @@ -434,6 +435,7 @@ public class ClusterReaderFactory { Set<String> deviceMeasurements, PartitionGroup partitionGroup, QueryContext context, boolean ascending) throws StorageEngineException { GroupByRequest request = new GroupByRequest(); + request.setRaftId(partitionGroup.getId()); if (timeFilter != null) { request.setTimeFilterBytes(SerializeUtils.serializeFilter(timeFilter)); } diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/server/heartbeat/HeartbeatThread.java b/cluster/src/main/java/org/apache/iotdb/cluster/server/heartbeat/HeartbeatThread.java index 0bfefee..2ccbaf3 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/server/heartbeat/HeartbeatThread.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/server/heartbeat/HeartbeatThread.java @@ -61,6 +61,7 @@ public class HeartbeatThread implements Runnable { HeartbeatThread(RaftMember localMember) { this.localMember = localMember; memberName = localMember.getName(); + request.setRaftId(localMember.getRaftGroupId()); } @Override @@ -202,6 +203,7 @@ public class HeartbeatThread implements Runnable { req.setRequireIdentifier(request.requireIdentifier); req.setTerm(request.term); req.setLeader(localMember.getThisNode()); + req.setRaftId(localMember.getRaftGroupId()); if (request.isSetHeader()) { req.setHeader(request.header); } diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/RaftMember.java b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/RaftMember.java index cbf0fe4..4b4d4e1 100644 --- a/cluster/src/main/java/org/apache/iotdb/cluster/server/member/RaftMember.java +++ b/cluster/src/main/java/org/apache/iotdb/cluster/server/member/RaftMember.java @@ -1421,6 +1421,7 @@ public abstract class RaftMember { AppendEntryRequest buildAppendEntryRequest(Log log, boolean serializeNow) { AppendEntryRequest request = new AppendEntryRequest(); + request.setRaftId(getRaftGroupId()); request.setTerm(term.get()); if (serializeNow) { request.setEntry(log.serialize()); diff --git a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java index 9847754..27af6b1 100644 --- a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java +++ b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBConfig.java @@ -113,7 +113,7 @@ public class IoTDBConfig { /** * whether to use Snappy compression before sending data through the network */ - private boolean rpcAdvancedCompressionEnable = true; + private boolean rpcAdvancedCompressionEnable = false; /** * Port which the JDBC server listens to. diff --git a/service-rpc/src/main/java/org/apache/iotdb/rpc/RpcTransportFactory.java b/service-rpc/src/main/java/org/apache/iotdb/rpc/RpcTransportFactory.java index 84f00df..57907de 100644 --- a/service-rpc/src/main/java/org/apache/iotdb/rpc/RpcTransportFactory.java +++ b/service-rpc/src/main/java/org/apache/iotdb/rpc/RpcTransportFactory.java @@ -27,7 +27,7 @@ import org.apache.thrift.transport.TTransportFactory; public class RpcTransportFactory extends TTransportFactory { // TODO: make it a config - public static boolean USE_SNAPPY = true; + public static boolean USE_SNAPPY = false; public static final RpcTransportFactory INSTANCE; static { INSTANCE = USE_SNAPPY ? diff --git a/thrift/src/main/thrift/cluster.thrift b/thrift/src/main/thrift/cluster.thrift index 549cd55..69a4dc5 100644 --- a/thrift/src/main/thrift/cluster.thrift +++ b/thrift/src/main/thrift/cluster.thrift @@ -41,7 +41,7 @@ struct HeartBeatRequest { // because a data server may play many data groups members, this is used to identify which // member should process the request or response. Only used in data group communication. 8: optional Node header - 9: optional int raftId + 9: required int raftId } // follower -> leader @@ -57,7 +57,6 @@ struct HeartBeatResponse { // because a data server may play many data groups members, this is used to identify which // member should process the request or response. Only used in data group communication. 7: optional Node header - 8: optional int raftId } // node -> node @@ -70,7 +69,7 @@ struct ElectionRequest { // because a data server may play many data groups members, this is used to identify which // member should process the request or response. Only used in data group communication. 5: optional Node header - 6: optional int raftId + 6: required int raftId 7: optional long dataLogLastIndex 8: optional long dataLogLastTerm } @@ -87,7 +86,7 @@ struct AppendEntryRequest { // because a data server may play many data groups members, this is used to identify which // member should process the request or response. Only used in data group communication. 7: optional Node header - 8: optional int raftId + 8: required int raftId } // leader -> follower @@ -102,7 +101,7 @@ struct AppendEntriesRequest { // because a data server may play many data groups members, this is used to identify which // member should process the request or response. Only used in data group communication. 7: optional Node header - 8: optional int raftId + 8: required int raftId } struct AddNodeResponse { @@ -148,14 +147,14 @@ struct SendSnapshotRequest { 1: required binary snapshotBytes // for data group 2: optional Node header - 3: optional int raftId + 3: required int raftId } struct PullSnapshotRequest { 1: required list<int> requiredSlots // for data group 2: optional Node header - 3: optional int raftId + 3: required int raftId // set to true if the previous holder has been removed from the cluster. // This will make the previous holder read-only so that different new // replicas can pull the same snapshot. @@ -169,13 +168,13 @@ struct PullSnapshotResp { struct ExecutNonQueryReq { 1: required binary planBytes 2: optional Node header - 3: optional int raftId + 3: required int raftId } struct PullSchemaRequest { 1: required list<string> prefixPaths 2: optional Node header - 3: optional int raftId + 3: required int raftId } struct PullSchemaResp {
