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 {

Reply via email to