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