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);