This is an automated email from the ASF dual-hosted git repository.

jiangtian pushed a commit to branch native_raft
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/native_raft by this push:
     new 45433f66d5 fix flow calculation refactor config ignoreStatemachine
45433f66d5 is described below

commit 45433f66d56592a19dfbeea24b39e731b1e9e173
Author: Tian Jiang <[email protected]>
AuthorDate: Thu May 11 11:07:06 2023 +0800

    fix flow calculation
    refactor config ignoreStatemachine
---
 .../apache/iotdb/consensus/EmptyStateMachine.java  |  2 ++
 .../log/dispatch/flowcontrol/FlowBalancer.java     | 12 +++++++--
 .../java/org/apache/iotdb/db/conf/IoTDBConfig.java | 10 ++++++++
 .../org/apache/iotdb/db/conf/IoTDBDescriptor.java  |  5 ++++
 .../db/consensus/DataRegionConsensusImpl.java      | 29 ++++++++++++++++------
 5 files changed, 48 insertions(+), 10 deletions(-)

diff --git 
a/consensus/src/test/java/org/apache/iotdb/consensus/EmptyStateMachine.java 
b/consensus/src/main/java/org/apache/iotdb/consensus/EmptyStateMachine.java
similarity index 98%
rename from 
consensus/src/test/java/org/apache/iotdb/consensus/EmptyStateMachine.java
rename to 
consensus/src/main/java/org/apache/iotdb/consensus/EmptyStateMachine.java
index 4580a865c2..daa8da958d 100644
--- a/consensus/src/test/java/org/apache/iotdb/consensus/EmptyStateMachine.java
+++ b/consensus/src/main/java/org/apache/iotdb/consensus/EmptyStateMachine.java
@@ -27,6 +27,8 @@ import java.io.File;
 
 public class EmptyStateMachine implements IStateMachine, 
IStateMachine.EventApi {
 
+  public EmptyStateMachine() {}
+
   @Override
   public void start() {}
 
diff --git 
a/consensus/src/main/java/org/apache/iotdb/consensus/natraft/protocol/log/dispatch/flowcontrol/FlowBalancer.java
 
b/consensus/src/main/java/org/apache/iotdb/consensus/natraft/protocol/log/dispatch/flowcontrol/FlowBalancer.java
index 985cb56985..3e2b408ccc 100644
--- 
a/consensus/src/main/java/org/apache/iotdb/consensus/natraft/protocol/log/dispatch/flowcontrol/FlowBalancer.java
+++ 
b/consensus/src/main/java/org/apache/iotdb/consensus/natraft/protocol/log/dispatch/flowcontrol/FlowBalancer.java
@@ -106,7 +106,7 @@ public class FlowBalancer {
         latestWindows.stream().mapToLong(w -> w.sum).sum()
             * 1.0
             / latestWindows.size()
-            * flowMonitorWindowInterval
+            * (flowMonitorWindowInterval / 1000.0)
             * overestimateFactor;
 
     Map<Peer, DispatcherGroup> dispatcherGroupMap = 
logDispatcher.getDispatcherGroupMap();
@@ -143,7 +143,6 @@ public class FlowBalancer {
 
   private void enterBurst(
       Map<Peer, Double> nodesRate, int nodeNum, double assumedFlow, List<Peer> 
followers) {
-    logger.info("{}: entering burst", member.getName());
     int followerNum = nodeNum - 1;
     int quorumFollowerNum = nodeNum / 2;
     double remainingFlow = maxFlow;
@@ -157,6 +156,15 @@ public class FlowBalancer {
       remainingFlow -= flowToQuorum;
     }
     double flowToRemaining = remainingFlow / (followerNum - quorumFollowerNum);
+    logger.info(
+        "{}: entering burst, quorum flow: {}, non-quorum flow: {}, max flow: 
{}, min flow: {}, assumed flow: {}, quorum max flow: {}",
+        member.getName(),
+        flowToQuorum,
+        flowToRemaining,
+        maxFlow,
+        minFlow,
+        assumedFlow,
+        quorumMaxFlow);
     if (flowToRemaining < minFlow) {
       flowToRemaining = minFlow;
     }
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 f0aa7876d1..f657ee635b 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
@@ -1116,6 +1116,8 @@ public class IoTDBConfig {
   /** A DynamicThread will not automatically exit unless its running time 
exceeds the value. */
   private long dynamicThreadMinRunningTimeNS = 10_000_000_000L;
 
+  private boolean ignoreStateMachine = false;
+
   IoTDBConfig() {}
 
   public float getUdfMemoryBudgetInMB() {
@@ -3854,4 +3856,12 @@ public class IoTDBConfig {
   public String getSortTmpDir() {
     return sortTmpDir;
   }
+
+  public boolean isIgnoreStateMachine() {
+    return ignoreStateMachine;
+  }
+
+  public void setIgnoreStateMachine(boolean ignoreStateMachine) {
+    this.ignoreStateMachine = ignoreStateMachine;
+  }
 }
diff --git a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java 
b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
index 4ad9e43c09..ea44771ed4 100644
--- a/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
+++ b/server/src/main/java/org/apache/iotdb/db/conf/IoTDBDescriptor.java
@@ -1045,6 +1045,11 @@ public class IoTDBDescriptor {
                 "dynamic_min_running_time_ns",
                 Long.toString(conf.getDynamicThreadMinRunningTimeNS()))));
 
+    conf.setIgnoreStateMachine(
+        Boolean.parseBoolean(
+            properties.getProperty(
+                "ignore_state_machine", 
String.valueOf(conf.isIgnoreStateMachine()))));
+
     // commons
     commonDescriptor.loadCommonProps(properties);
     commonDescriptor.initCommonConfigDir(conf.getSystemDir());
diff --git 
a/server/src/main/java/org/apache/iotdb/db/consensus/DataRegionConsensusImpl.java
 
b/server/src/main/java/org/apache/iotdb/db/consensus/DataRegionConsensusImpl.java
index 5684202a0c..69823a25e8 100644
--- 
a/server/src/main/java/org/apache/iotdb/db/consensus/DataRegionConsensusImpl.java
+++ 
b/server/src/main/java/org/apache/iotdb/db/consensus/DataRegionConsensusImpl.java
@@ -23,11 +23,20 @@ import org.apache.iotdb.common.rpc.thrift.TEndPoint;
 import org.apache.iotdb.commons.consensus.ConsensusGroupId;
 import org.apache.iotdb.commons.consensus.DataRegionId;
 import org.apache.iotdb.consensus.ConsensusFactory;
+import org.apache.iotdb.consensus.EmptyStateMachine;
 import org.apache.iotdb.consensus.IConsensus;
+import org.apache.iotdb.consensus.IStateMachine;
 import org.apache.iotdb.consensus.config.ConsensusConfig;
 import org.apache.iotdb.consensus.config.IoTConsensusConfig;
 import org.apache.iotdb.consensus.config.IoTConsensusConfig.RPC;
+import org.apache.iotdb.consensus.config.IoTConsensusConfig.Replication;
 import org.apache.iotdb.consensus.config.RatisConfig;
+import org.apache.iotdb.consensus.config.RatisConfig.Client;
+import org.apache.iotdb.consensus.config.RatisConfig.Grpc;
+import org.apache.iotdb.consensus.config.RatisConfig.Impl;
+import org.apache.iotdb.consensus.config.RatisConfig.LeaderLogAppender;
+import org.apache.iotdb.consensus.config.RatisConfig.Log;
+import org.apache.iotdb.consensus.config.RatisConfig.Rpc;
 import org.apache.iotdb.consensus.config.RatisConfig.Snapshot;
 import org.apache.iotdb.db.conf.IoTDBConfig;
 import org.apache.iotdb.db.conf.IoTDBDescriptor;
@@ -93,7 +102,7 @@ public class DataRegionConsensusImpl {
                                       
.setMaxClientNumForEachNode(conf.getMaxClientNumForEachNode())
                                       .build())
                               .setReplication(
-                                  IoTConsensusConfig.Replication.newBuilder()
+                                  Replication.newBuilder()
                                       
.setWalThrottleThreshold(conf.getThrottleThreshold())
                                       .setAllocateMemoryForConsensus(
                                           conf.getAllocateMemoryForConsensus())
@@ -111,7 +120,7 @@ public class DataRegionConsensusImpl {
                                           
conf.getDataRatisConsensusSnapshotTriggerThreshold())
                                       .build())
                               .setLog(
-                                  RatisConfig.Log.newBuilder()
+                                  Log.newBuilder()
                                       .setUnsafeFlushEnabled(
                                           
conf.isDataRatisConsensusLogUnsafeFlushEnable())
                                       .setSegmentSizeMax(
@@ -121,13 +130,13 @@ public class DataRegionConsensusImpl {
                                           
conf.getDataRatisConsensusPreserveWhenPurge())
                                       .build())
                               .setGrpc(
-                                  RatisConfig.Grpc.newBuilder()
+                                  Grpc.newBuilder()
                                       .setFlowControlWindow(
                                           SizeInBytes.valueOf(
                                               
conf.getDataRatisConsensusGrpcFlowControlWindow()))
                                       .build())
                               .setRpc(
-                                  RatisConfig.Rpc.newBuilder()
+                                  Rpc.newBuilder()
                                       .setTimeoutMin(
                                           TimeDuration.valueOf(
                                               conf
@@ -152,7 +161,7 @@ public class DataRegionConsensusImpl {
                                               TimeUnit.MILLISECONDS))
                                       .build())
                               .setClient(
-                                  RatisConfig.Client.newBuilder()
+                                  Client.newBuilder()
                                       .setClientRequestTimeoutMillis(
                                           
conf.getDataRatisConsensusRequestTimeoutMs())
                                       .setClientMaxRetryAttempt(
@@ -166,11 +175,11 @@ public class DataRegionConsensusImpl {
                                       
.setMaxClientNumForEachNode(conf.getMaxClientNumForEachNode())
                                       .build())
                               .setImpl(
-                                  RatisConfig.Impl.newBuilder()
+                                  Impl.newBuilder()
                                       
.setTriggerSnapshotFileSize(conf.getDataRatisLogMax())
                                       .build())
                               .setLeaderLogAppender(
-                                  RatisConfig.LeaderLogAppender.newBuilder()
+                                  LeaderLogAppender.newBuilder()
                                       .setBufferByteLimit(
                                           
conf.getDataRatisConsensusLogAppenderBufferSizeMax())
                                       .build())
@@ -188,7 +197,11 @@ public class DataRegionConsensusImpl {
     return INSTANCE;
   }
 
-  private static DataRegionStateMachine 
createDataRegionStateMachine(ConsensusGroupId gid) {
+  private static IStateMachine createDataRegionStateMachine(ConsensusGroupId 
gid) {
+    if (conf.isIgnoreStateMachine()) {
+      return new EmptyStateMachine();
+    }
+
     DataRegion dataRegion = 
StorageEngine.getInstance().getDataRegion((DataRegionId) gid);
     if 
(ConsensusFactory.IOT_CONSENSUS.equals(conf.getDataRegionConsensusProtocolClass()))
 {
       return new IoTConsensusDataRegionStateMachine(dataRegion);

Reply via email to