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

Reply via email to