This is an automated email from the ASF dual-hosted git repository. xingtanzjr pushed a commit to branch multiLeader_out_of_order in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit 2fc15d31e5c324d8e46e69824d26a4dc1e783849 Author: stormbroken <[email protected]> AuthorDate: Mon Jul 25 09:39:23 2022 +0800 add some log. --- .../common/request/IndexedConsensusRequest.java | 20 ++------------------ .../consensus/multileader/MultiLeaderServerImpl.java | 4 ++-- .../multileader/logdispatcher/LogDispatcher.java | 11 +++++++++-- 3 files changed, 13 insertions(+), 22 deletions(-) diff --git a/consensus/src/main/java/org/apache/iotdb/consensus/common/request/IndexedConsensusRequest.java b/consensus/src/main/java/org/apache/iotdb/consensus/common/request/IndexedConsensusRequest.java index de7ac07250..de3aca433b 100644 --- a/consensus/src/main/java/org/apache/iotdb/consensus/common/request/IndexedConsensusRequest.java +++ b/consensus/src/main/java/org/apache/iotdb/consensus/common/request/IndexedConsensusRequest.java @@ -19,8 +19,6 @@ package org.apache.iotdb.consensus.common.request; -import org.apache.iotdb.consensus.multileader.wal.ConsensusReqReader; - import java.nio.ByteBuffer; import java.util.List; import java.util.Objects; @@ -31,18 +29,10 @@ public class IndexedConsensusRequest implements IConsensusRequest { /** we do not need to serialize these two fields as they are useless in other nodes. */ private final long searchIndex; - private final long safelyDeletedSearchIndex; - private final List<IConsensusRequest> requests; public IndexedConsensusRequest(long searchIndex, List<IConsensusRequest> requests) { - this(searchIndex, ConsensusReqReader.DEFAULT_SAFELY_DELETED_SEARCH_INDEX, requests); - } - - public IndexedConsensusRequest( - long searchIndex, long safelyDeletedSearchIndex, List<IConsensusRequest> requests) { this.searchIndex = searchIndex; - this.safelyDeletedSearchIndex = safelyDeletedSearchIndex; this.requests = requests; } @@ -59,10 +49,6 @@ public class IndexedConsensusRequest implements IConsensusRequest { return searchIndex; } - public long getSafelyDeletedSearchIndex() { - return safelyDeletedSearchIndex; - } - @Override public boolean equals(Object o) { if (this == o) { @@ -72,13 +58,11 @@ public class IndexedConsensusRequest implements IConsensusRequest { return false; } IndexedConsensusRequest that = (IndexedConsensusRequest) o; - return searchIndex == that.searchIndex - && safelyDeletedSearchIndex == that.safelyDeletedSearchIndex - && requests.equals(that.requests); + return searchIndex == that.searchIndex && requests.equals(that.requests); } @Override public int hashCode() { - return Objects.hash(searchIndex, safelyDeletedSearchIndex, requests); + return Objects.hash(searchIndex, requests); } } diff --git a/consensus/src/main/java/org/apache/iotdb/consensus/multileader/MultiLeaderServerImpl.java b/consensus/src/main/java/org/apache/iotdb/consensus/multileader/MultiLeaderServerImpl.java index a0439bf8c9..d74d1613a7 100644 --- a/consensus/src/main/java/org/apache/iotdb/consensus/multileader/MultiLeaderServerImpl.java +++ b/consensus/src/main/java/org/apache/iotdb/consensus/multileader/MultiLeaderServerImpl.java @@ -114,9 +114,9 @@ public class MultiLeaderServerImpl { buildIndexedConsensusRequestForLocalRequest(request); if (indexedConsensusRequest.getSearchIndex() % 1000 == 0) { logger.info( - "DataRegion[{}]: index after build: safeIndex: {}, searchIndex: {}", + "DataRegion[{}]: index after build: safeIndex:{}, searchIndex: {}", thisNode.getGroupId(), - indexedConsensusRequest.getSafelyDeletedSearchIndex(), + getCurrentSafelyDeletedSearchIndex(), indexedConsensusRequest.getSearchIndex()); } // TODO wal and memtable diff --git a/consensus/src/main/java/org/apache/iotdb/consensus/multileader/logdispatcher/LogDispatcher.java b/consensus/src/main/java/org/apache/iotdb/consensus/multileader/logdispatcher/LogDispatcher.java index a9329d7d4d..a27357d304 100644 --- a/consensus/src/main/java/org/apache/iotdb/consensus/multileader/logdispatcher/LogDispatcher.java +++ b/consensus/src/main/java/org/apache/iotdb/consensus/multileader/logdispatcher/LogDispatcher.java @@ -34,11 +34,13 @@ import org.apache.iotdb.consensus.multileader.thrift.TSyncLogReq; import org.apache.iotdb.consensus.multileader.wal.ConsensusReqReader; import org.apache.iotdb.consensus.multileader.wal.GetConsensusReqReaderPlan; import org.apache.iotdb.consensus.ratis.Utils; + import org.apache.thrift.TException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.io.IOException; +import java.util.ArrayList; import java.util.Iterator; import java.util.LinkedList; import java.util.List; @@ -218,9 +220,9 @@ public class LogDispatcher { public PendingBatch getBatch() { PendingBatch batch; - List<TLogBatch> logBatches = new LinkedList<>(); + List<TLogBatch> logBatches = new ArrayList<>(); long startIndex = syncStatus.getNextSendingIndex(); - logger.debug("get batch. startIndex: {}", startIndex); + logger.debug("[GetBatch] startIndex: {}", startIndex); long endIndex; if (bufferedRequest.size() <= config.getReplication().getMaxRequestPerBatch()) { // Use drainTo instead of poll to reduce lock overhead @@ -300,6 +302,11 @@ public class LogDispatcher { AsyncMultiLeaderServiceClient client = clientManager.borrowClient(peer.getEndpoint()); TSyncLogReq req = new TSyncLogReq(peer.getGroupId().convertToTConsensusGroupId(), batch.getBatches()); + logger.info( + "Send Batch[startIndex:{}, endIndex:{}] to ConsensusGroup:{}", + batch.getStartIndex(), + batch.getEndIndex(), + peer.getGroupId().convertToTConsensusGroupId()); client.syncLog(req, handler); } catch (IOException | TException e) { logger.error("Can not sync logs to peer {} because", peer, e);
