This is an automated email from the ASF dual-hosted git repository.
jiangtian pushed a commit to branch expr
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/expr by this push:
new 2dba555 add configs
2dba555 is described below
commit 2dba555b6fb3d02a3dfa81bb3a4fb64a483d83df
Author: jt <[email protected]>
AuthorDate: Wed Dec 1 14:47:36 2021 +0800
add configs
---
.../apache/iotdb/cluster/config/ClusterConfig.java | 20 +++
.../iotdb/cluster/server/member/RaftMember.java | 150 ++++++++++++++-------
2 files changed, 120 insertions(+), 50 deletions(-)
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 b677867..e746b3a 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
@@ -184,6 +184,10 @@ public class ClusterConfig {
private boolean useAsyncSequencing = true;
+ private boolean useFollowerSlidingWindow = true;
+
+ private boolean enableWeakAcceptance = true;
+
/**
* create a clusterConfig class. The internalIP will be set according to the
server's hostname. If
* there is something error for getting the ip of the hostname, then set the
internalIp as
@@ -550,4 +554,20 @@ public class ClusterConfig {
public void setUseAsyncSequencing(boolean useAsyncSequencing) {
this.useAsyncSequencing = useAsyncSequencing;
}
+
+ public boolean isUseFollowerSlidingWindow() {
+ return useFollowerSlidingWindow;
+ }
+
+ public void setUseFollowerSlidingWindow(boolean useFollowerSlidingWindow) {
+ this.useFollowerSlidingWindow = useFollowerSlidingWindow;
+ }
+
+ public boolean isEnableWeakAcceptance() {
+ return enableWeakAcceptance;
+ }
+
+ public void setEnableWeakAcceptance(boolean enableWeakAcceptance) {
+ this.enableWeakAcceptance = enableWeakAcceptance;
+ }
}
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 c822ad5..4d4587d 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
@@ -44,6 +44,7 @@ import org.apache.iotdb.cluster.log.VotingLogList;
import org.apache.iotdb.cluster.log.appender.BlockingLogAppender;
import org.apache.iotdb.cluster.log.appender.LogAppender;
import org.apache.iotdb.cluster.log.appender.LogAppenderFactory;
+import org.apache.iotdb.cluster.log.appender.SlidingWindowLogAppender;
import org.apache.iotdb.cluster.log.catchup.CatchUpTask;
import org.apache.iotdb.cluster.log.logtypes.PhysicalPlanLog;
import org.apache.iotdb.cluster.log.manage.RaftLogManager;
@@ -129,16 +130,22 @@ import java.util.concurrent.atomic.AtomicReference;
import static
org.apache.iotdb.cluster.config.ClusterConstant.THREAD_POLL_WAIT_TERMINATION_TIME_S;
/**
- * RaftMember process the common raft logic like leader election, log
appending, catch-up and so on.
+ * RaftMember process the common raft logic like leader election, log
appending, catch-up and so
+ * on.
*/
@SuppressWarnings("java:S3077") // reference volatile is enough
public abstract class RaftMember implements RaftMemberMBean {
+
private static final Logger logger =
LoggerFactory.getLogger(RaftMember.class);
public static boolean USE_LOG_DISPATCHER = false;
- public static boolean USE_INDIRECT_LOG_DISPATCHER = false;
- public static boolean ENABLE_WEAK_ACCEPTANCE = true;
-
- private static final LogAppenderFactory APPENDER_FACTORY = new
BlockingLogAppender.Factory();
+ private static final boolean USE_INDIRECT_LOG_DISPATCHER =
+ ClusterDescriptor.getInstance().getConfig().isUseIndirectBroadcasting();
+ private static final boolean ENABLE_WEAK_ACCEPTANCE =
ClusterDescriptor.getInstance().getConfig()
+ .isEnableWeakAcceptance();
+
+ private static final LogAppenderFactory APPENDER_FACTORY =
+ ClusterDescriptor.getInstance().getConfig().isUseFollowerSlidingWindow()
?
+ new SlidingWindowLogAppender.Factory() : new
BlockingLogAppender.Factory();
protected static final LogSequencerFactory SEQUENCER_FACTORY =
ClusterDescriptor.getInstance().getConfig().isUseAsyncSequencing()
? new Factory()
@@ -162,22 +169,32 @@ public abstract class RaftMember implements
RaftMemberMBean {
* on this may be woken.
*/
private final Object waitLeaderCondition = new Object();
- /** the lock is to make sure that only one thread can apply snapshot at the
same time */
+ /**
+ * the lock is to make sure that only one thread can apply snapshot at the
same time
+ */
private final Object snapshotApplyLock = new Object();
private final Object heartBeatWaitObject = new Object();
protected Node thisNode = ClusterIoTDB.getInstance().getThisNode();
- /** the nodes that belong to the same raft group as thisNode. */
+ /**
+ * the nodes that belong to the same raft group as thisNode.
+ */
protected PartitionGroup allNodes;
ClusterConfig config = ClusterDescriptor.getInstance().getConfig();
- /** the name of the member, to distinguish several members in the logs. */
+ /**
+ * the name of the member, to distinguish several members in the logs.
+ */
String name;
- /** to choose nodes to send request of joining cluster randomly. */
+ /**
+ * to choose nodes to send request of joining cluster randomly.
+ */
Random random = new Random();
- /** when the node is a leader, this map is used to track log progress of
each follower. */
+ /**
+ * when the node is a leader, this map is used to track log progress of each
follower.
+ */
Map<Node, Peer> peerMap;
/**
* the current term of the node, this object also works as lock of some
transactions of the member
@@ -199,7 +216,9 @@ public abstract class RaftMember implements RaftMemberMBean
{
*/
volatile long lastHeartbeatReceivedTime;
- /** the raft logs are all stored and maintained in the log manager */
+ /**
+ * the raft logs are all stored and maintained in the log manager
+ */
protected RaftLogManager logManager;
/**
@@ -218,7 +237,9 @@ public abstract class RaftMember implements RaftMemberMBean
{
* member by comparing it with the current last log index.
*/
long lastReportedLogIndex;
- /** the thread pool that runs catch-up tasks */
+ /**
+ * the thread pool that runs catch-up tasks
+ */
private ExecutorService catchUpService;
/**
* lastCatchUpResponseTime records when is the latest response of each
node's catch-up. There
@@ -249,24 +270,32 @@ public abstract class RaftMember implements
RaftMemberMBean {
* one slow node.
*/
private ExecutorService serialToParallelPool;
- /** a thread pool that is used to do commit log tasks asynchronous in
heartbeat thread */
+ /**
+ * a thread pool that is used to do commit log tasks asynchronous in
heartbeat thread
+ */
private ExecutorService commitLogPool;
/**
* logDispatcher buff the logs orderly according to their log indexes and
send them sequentially,
- * which avoids the followers receiving out-of-order logs, forcing them to
wait for previous logs.
+ * which avoids the followers receiving out-of-order logs, forcing them to
wait for previous
+ * logs.
*/
private volatile LogDispatcher logDispatcher;
- /** If this node can not be the leader, this parameter will be set true. */
+ /**
+ * If this node can not be the leader, this parameter will be set true.
+ */
private volatile boolean skipElection = false;
/**
- * localExecutor is used to directly execute plans like load configuration
in the underlying IoTDB
+ * localExecutor is used to directly execute plans like load configuration
in the underlying
+ * IoTDB
*/
protected PlanExecutor localExecutor;
- /** (logIndex, logTerm) -> append handler */
+ /**
+ * (logIndex, logTerm) -> append handler
+ */
protected Map<Pair<Long, Long>, AppendNodeEntryHandler> sentLogHandlers =
new ConcurrentHashMap<>();
@@ -276,7 +305,8 @@ public abstract class RaftMember implements RaftMemberMBean
{
private volatile LogAppender logAppender;
- protected RaftMember() {}
+ protected RaftMember() {
+ }
protected RaftMember(String name, ClientManager clientManager) {
this.name = name;
@@ -632,7 +662,9 @@ public abstract class RaftMember implements RaftMemberMBean
{
}
}
- /** Similar to appendEntry, while the incoming load is batch of logs instead
of a single log. */
+ /**
+ * Similar to appendEntry, while the incoming load is batch of logs instead
of a single log.
+ */
public AppendEntryResult appendEntries(AppendEntriesRequest request)
throws UnknownLogTypeException {
logger.debug("{} received an AppendEntriesRequest", name);
@@ -786,16 +818,22 @@ public abstract class RaftMember implements
RaftMemberMBean {
return lastCatchUpResponseTime;
}
- /** Sub-classes will add their own process of HeartBeatResponse in this
method. */
- public void processValidHeartbeatResp(HeartBeatResponse response, Node
receiver) {}
+ /**
+ * Sub-classes will add their own process of HeartBeatResponse in this
method.
+ */
+ public void processValidHeartbeatResp(HeartBeatResponse response, Node
receiver) {
+ }
- /** The actions performed when the node wins in an election (becoming a
leader). */
- public void onElectionWins() {}
+ /**
+ * The actions performed when the node wins in an election (becoming a
leader).
+ */
+ public void onElectionWins() {
+ }
/**
* Update the followers' log by sending logs whose index >=
followerLastMatchedLogIndex to the
- * follower. If some of the required logs are removed, also send the
snapshot. <br>
- * notice that if a part of data is in the snapshot, then it is not in the
logs.
+ * follower. If some of the required logs are removed, also send the
snapshot. <br> notice that if
+ * a part of data is in the snapshot, then it is not in the logs.
*/
public void catchUp(Node follower, long lastLogIdx) {
// for one follower, there is at most one ongoing catch-up, so the same
data will not be sent
@@ -889,7 +927,9 @@ public abstract class RaftMember implements RaftMemberMBean
{
"%s:%s=%s", "org.apache.iotdb.cluster.service",
IoTDBConstant.JMX_TYPE, "Engine");
}
- /** call back after syncLeader */
+ /**
+ * call back after syncLeader
+ */
public interface CheckConsistency {
/**
@@ -898,7 +938,7 @@ public abstract class RaftMember implements RaftMemberMBean
{
* @param leaderCommitId leader commit id
* @param localAppliedId local applied id
* @throws CheckConsistencyException maybe throw
CheckConsistencyException, which is defined in
- * implements.
+ * implements.
*/
void postCheckConsistency(long leaderCommitId, long localAppliedId)
throws CheckConsistencyException;
@@ -907,8 +947,7 @@ public abstract class RaftMember implements RaftMemberMBean
{
public static class MidCheckConsistency implements CheckConsistency {
/**
- * if leaderCommitId - localAppliedId > MaxReadLogLag, will throw
- * CHECK_MID_CONSISTENCY_EXCEPTION
+ * if leaderCommitId - localAppliedId > MaxReadLogLag, will throw
CHECK_MID_CONSISTENCY_EXCEPTION
*
* @param leaderCommitId leader commit id
* @param localAppliedId local applied id
@@ -920,7 +959,7 @@ public abstract class RaftMember implements RaftMemberMBean
{
if (leaderCommitId == Long.MAX_VALUE
|| leaderCommitId == Long.MIN_VALUE
|| leaderCommitId - localAppliedId
- >
ClusterDescriptor.getInstance().getConfig().getMaxReadLogLag()) {
+ > ClusterDescriptor.getInstance().getConfig().getMaxReadLogLag()) {
throw CheckConsistencyException.CHECK_MID_CONSISTENCY_EXCEPTION;
}
}
@@ -953,7 +992,7 @@ public abstract class RaftMember implements RaftMemberMBean
{
* @param checkConsistency check after syncleader
* @return true if the node has caught up, false otherwise
* @throws CheckConsistencyException if leaderCommitId bigger than
localAppliedId a threshold
- * value after timeout
+ * value after timeout
*/
public boolean syncLeader(CheckConsistency checkConsistency) throws
CheckConsistencyException {
if (character == NodeCharacter.LEADER) {
@@ -972,7 +1011,9 @@ public abstract class RaftMember implements
RaftMemberMBean {
return waitUntilCatchUp(checkConsistency);
}
- /** Wait until the leader of this node becomes known or time out. */
+ /**
+ * Wait until the leader of this node becomes known or time out.
+ */
public void waitLeader() {
long startTime = System.currentTimeMillis();
while (leader.get() == null ||
ClusterConstant.EMPTY_NODE.equals(leader.get())) {
@@ -999,7 +1040,7 @@ public abstract class RaftMember implements
RaftMemberMBean {
*
* @return true if this node has caught up before timeout, false otherwise
* @throws CheckConsistencyException if leaderCommitId bigger than
localAppliedId a threshold
- * value after timeout
+ * value after timeout
*/
protected boolean waitUntilCatchUp(CheckConsistency checkConsistency)
throws CheckConsistencyException {
@@ -1032,7 +1073,7 @@ public abstract class RaftMember implements
RaftMemberMBean {
* sync local applyId to leader commitId
*
* @param leaderCommitId leader commit id
- * @param fastFail if enable, when log differ too much, return false
directly.
+ * @param fastFail if enable, when log differ too much, return false
directly.
* @return true if leaderCommitId <= localAppliedId
*/
public boolean syncLocalApply(long leaderCommitId, boolean fastFail) {
@@ -1085,7 +1126,7 @@ public abstract class RaftMember implements
RaftMemberMBean {
* call this method. Will commit the log locally and send it to followers
*
* @return OK if over half of the followers accept the log or null if the
leadership is lost
- * during the appending
+ * during the appending
*/
public TSStatus processPlanLocally(PhysicalPlan plan) {
if (USE_LOG_DISPATCHER) {
@@ -1322,7 +1363,9 @@ public abstract class RaftMember implements
RaftMemberMBean {
return peerMap;
}
- /** @return true if there is a log whose index is "index" and term is
"term", false otherwise */
+ /**
+ * @return true if there is a log whose index is "index" and term is "term",
false otherwise
+ */
public boolean matchLog(long index, long term) {
boolean matched = logManager.matchTerm(term, index);
logger.debug("Log {}-{} matched: {}", index, term, matched);
@@ -1341,15 +1384,18 @@ public abstract class RaftMember implements
RaftMemberMBean {
return syncLock;
}
- /** Sub-classes will add their own process of HeartBeatRequest in this
method. */
- void processValidHeartbeatReq(HeartBeatRequest request, HeartBeatResponse
response) {}
+ /**
+ * Sub-classes will add their own process of HeartBeatRequest in this method.
+ */
+ void processValidHeartbeatReq(HeartBeatRequest request, HeartBeatResponse
response) {
+ }
/**
* Verify the validity of an ElectionRequest, and make itself a follower of
the elector if the
* request is valid.
*
* @return Response.RESPONSE_AGREE if the elector is valid or the local term
if the elector has a
- * smaller term or Response.RESPONSE_LOG_MISMATCH if the elector has
older logs.
+ * smaller term or Response.RESPONSE_LOG_MISMATCH if the elector has older
logs.
*/
long checkElectorLogProgress(ElectionRequest electionRequest) {
@@ -1393,7 +1439,7 @@ public abstract class RaftMember implements
RaftMemberMBean {
* lastLogIndex is smaller than the voter's Otherwise accept the election.
*
* @return Response.RESPONSE_AGREE if the elector is valid or the local term
if the elector has a
- * smaller term or Response.RESPONSE_LOG_MISMATCH if the elector has
older logs.
+ * smaller term or Response.RESPONSE_LOG_MISMATCH if the elector has older
logs.
*/
long checkLogProgress(long lastLogIndex, long lastLogTerm) {
long response;
@@ -1410,10 +1456,10 @@ public abstract class RaftMember implements
RaftMemberMBean {
/**
* Forward a non-query plan to a node using the default client.
*
- * @param plan a non-query plan
- * @param node cannot be the local node
+ * @param plan a non-query plan
+ * @param node cannot be the local node
* @param header must be set for data group communication, set to null for
meta group
- * communication
+ * communication
* @return a TSStatus indicating if the forwarding is successful.
*/
public TSStatus forwardPlan(PhysicalPlan plan, Node node, RaftNode header) {
@@ -1444,7 +1490,7 @@ public abstract class RaftMember implements
RaftMemberMBean {
/**
* Forward a non-query plan to "receiver" using "client".
*
- * @param plan a non-query plan
+ * @param plan a non-query plan
* @param header to determine which DataGroupMember of "receiver" will
process the request.
* @return a TSStatus indicating if the forwarding is successful.
*/
@@ -1526,7 +1572,7 @@ public abstract class RaftMember implements
RaftMemberMBean {
* Get an asynchronous thrift client of the given node.
*
* @return an asynchronous thrift client or null if the caller tries to
connect the local node or
- * the node cannot be reached.
+ * the node cannot be reached.
*/
public AsyncClient getAsyncClient(Node node) {
try {
@@ -1654,8 +1700,8 @@ public abstract class RaftMember implements
RaftMemberMBean {
long alreadyWait = 0;
while (stronglyAcceptedNodeNum < quorumSize
&& (!ENABLE_WEAK_ACCEPTANCE
- || (totalAccepted < allNodes.size() - 1)
- || votingLogList.size() > config.getMaxNumOfLogsInMem())
+ || (totalAccepted < allNodes.size() - 1)
+ || votingLogList.size() > config.getMaxNumOfLogsInMem())
&& alreadyWait < ClusterConstant.getWriteOperationTimeoutMS()
&& !log.getStronglyAcceptedNodeIds().contains(Integer.MAX_VALUE)) {
try {
@@ -1817,7 +1863,7 @@ public abstract class RaftMember implements
RaftMemberMBean {
* heartbeat timer.
*
* @param fromLeader true if the request is from a leader, false if the
request is from an
- * elector.
+ * elector.
*/
public void stepDown(long newTerm, boolean fromLeader) {
synchronized (term) {
@@ -1849,7 +1895,9 @@ public abstract class RaftMember implements
RaftMemberMBean {
this.thisNode = thisNode;
}
- /** @return the header of the data raft group or null if this is in a meta
group. */
+ /**
+ * @return the header of the data raft group or null if this is in a meta
group.
+ */
public RaftNode getHeader() {
return null;
}
@@ -2017,7 +2065,9 @@ public abstract class RaftMember implements
RaftMemberMBean {
log, node, leaderShipStale, newLeaderTerm, request, quorumSize,
Collections.emptyList());
}
- /** Send "log" to "node". */
+ /**
+ * Send "log" to "node".
+ */
public void sendLogToFollower(
VotingLog log,
Node node,