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

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

commit aea0498d5d94ccf56c260da1294df913f7b991b9
Author: jt <[email protected]>
AuthorDate: Thu Apr 22 09:48:31 2021 +0800

    add logs to detect log application stuck
---
 cluster/src/main/java/org/apache/iotdb/cluster/log/Log.java  | 11 +++++++++++
 .../iotdb/cluster/log/applier/AsyncDataLogApplier.java       |  2 +-
 .../org/apache/iotdb/cluster/log/applier/DataLogApplier.java |  1 +
 .../org/apache/iotdb/cluster/log/manage/RaftLogManager.java  | 12 ++++++++++++
 4 files changed, 25 insertions(+), 1 deletion(-)

diff --git a/cluster/src/main/java/org/apache/iotdb/cluster/log/Log.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/log/Log.java
index e70c326..5498962 100644
--- a/cluster/src/main/java/org/apache/iotdb/cluster/log/Log.java
+++ b/cluster/src/main/java/org/apache/iotdb/cluster/log/Log.java
@@ -47,6 +47,9 @@ public abstract class Log implements Comparable<Log> {
 
   private int byteSize = 0;
 
+  // the thread that is applying this log, for debugging stuck applications
+  private Thread applyingThread;
+
   public abstract ByteBuffer serialize();
 
   public abstract void deserialize(ByteBuffer buffer);
@@ -142,4 +145,12 @@ public abstract class Log implements Comparable<Log> {
   public void setByteSize(int byteSize) {
     this.byteSize = byteSize;
   }
+
+  public Thread getApplyingThread() {
+    return applyingThread;
+  }
+
+  public void setApplyingThread(Thread applyingThread) {
+    this.applyingThread = applyingThread;
+  }
 }
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/log/applier/AsyncDataLogApplier.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/log/applier/AsyncDataLogApplier.java
index 65a5b2d..3b960f6 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/log/applier/AsyncDataLogApplier.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/log/applier/AsyncDataLogApplier.java
@@ -87,7 +87,7 @@ public class AsyncDataLogApplier implements LogApplier {
   // synchronized: when a log is draining consumers, avoid other threads 
adding more logs so that
   // the consumers will never be drained
   public synchronized void apply(Log log) {
-
+    log.setApplyingThread(Thread.currentThread());
     PartialPath logKey;
     try {
       logKey = getLogKey(log);
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/log/applier/DataLogApplier.java
 
b/cluster/src/main/java/org/apache/iotdb/cluster/log/applier/DataLogApplier.java
index 7d11d3e..1ccd94a 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/log/applier/DataLogApplier.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/log/applier/DataLogApplier.java
@@ -61,6 +61,7 @@ public class DataLogApplier extends BaseApplier {
 
   @Override
   public void apply(Log log) {
+    log.setApplyingThread(Thread.currentThread());
     logger.debug("DataMember [{}] start applying Log {}", 
dataGroupMember.getName(), log);
 
     try {
diff --git 
a/cluster/src/main/java/org/apache/iotdb/cluster/log/manage/RaftLogManager.java 
b/cluster/src/main/java/org/apache/iotdb/cluster/log/manage/RaftLogManager.java
index 446eefc..d818990 100644
--- 
a/cluster/src/main/java/org/apache/iotdb/cluster/log/manage/RaftLogManager.java
+++ 
b/cluster/src/main/java/org/apache/iotdb/cluster/log/manage/RaftLogManager.java
@@ -19,6 +19,7 @@
 
 package org.apache.iotdb.cluster.log.manage;
 
+import java.util.Arrays;
 import org.apache.iotdb.cluster.config.ClusterDescriptor;
 import org.apache.iotdb.cluster.exception.EntryCompactedException;
 import org.apache.iotdb.cluster.exception.EntryUnavailableException;
@@ -943,10 +944,21 @@ public abstract class RaftLogManager {
             nextToCheckIndex);
         return;
       }
+      boolean stuckLogPrinted = false;
+      long waitedTime;
+      long waitStartTime = System.currentTimeMillis();
       synchronized (log) {
         while (!log.isApplied() && maxHaveAppliedCommitIndex < 
log.getCurrLogIndex()) {
           // wait until the log is applied or a newer snapshot is installed
           log.wait(5);
+          waitedTime = System.currentTimeMillis() - waitStartTime;
+          if (!stuckLogPrinted && waitedTime > 60_000L) {
+            Thread applyingThread = log.getApplyingThread();
+            logger.error("The application of log {} is stuck for over 1 
minutes, processing "
+                + "thread: {}, thread traces: {}", log, applyingThread,
+                applyingThread != null ? 
Arrays.toString(applyingThread.getStackTrace()) : "null");
+            stuckLogPrinted = true;
+          }
         }
       }
       synchronized (changeApplyCommitIndexCond) {

Reply via email to