This is an automated email from the ASF dual-hosted git repository.
jiangtian pushed a commit to branch fix_ratis_restart
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/fix_ratis_restart by this push:
new 51d409f31c9 Fix that stopping ratis may get stuck due to closed
FlushManager
51d409f31c9 is described below
commit 51d409f31c9782b6caf8d23af2fef6bc5d7b50b0
Author: Tian Jiang <[email protected]>
AuthorDate: Tue Jul 22 11:40:12 2025 +0800
Fix that stopping ratis may get stuck due to closed FlushManager
---
.../apache/iotdb/db/service/DataNodeShutdownHook.java | 2 +-
.../storageengine/rescon/memory/AbstractPoolManager.java | 16 +++++++++++++++-
2 files changed, 16 insertions(+), 2 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/DataNodeShutdownHook.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/DataNodeShutdownHook.java
index e9b125818fa..1b6eeb82781 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/DataNodeShutdownHook.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/service/DataNodeShutdownHook.java
@@ -183,7 +183,7 @@ public class DataNodeShutdownHook extends Thread {
.forEach(
id -> {
try {
- DataRegionConsensusImpl.getInstance().triggerSnapshot(id,
false);
+ DataRegionConsensusImpl.getInstance().triggerSnapshot(id,
true);
} catch (ConsensusException e) {
logger.warn(
"Something wrong happened while calling consensus layer's "
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/rescon/memory/AbstractPoolManager.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/rescon/memory/AbstractPoolManager.java
index 405503e872e..c57459d2491 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/rescon/memory/AbstractPoolManager.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/rescon/memory/AbstractPoolManager.java
@@ -22,6 +22,7 @@ package org.apache.iotdb.db.storageengine.rescon.memory;
import org.slf4j.Logger;
import java.util.concurrent.Callable;
+import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Future;
import java.util.concurrent.ThreadPoolExecutor;
@@ -53,10 +54,23 @@ public abstract class AbstractPoolManager {
}
public synchronized Future<?> submit(Runnable task) {
+ if (pool == null) {
+ return CompletableFuture.runAsync(task);
+ }
return pool.submit(task);
}
public synchronized <T> Future<T> submit(Callable<T> task) {
+ if (pool == null) {
+ return CompletableFuture.supplyAsync(
+ () -> {
+ try {
+ return task.call();
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+ });
+ }
return pool.submit(task);
}
@@ -91,7 +105,7 @@ public abstract class AbstractPoolManager {
public abstract void start();
- public void stop() {
+ public synchronized void stop() {
if (pool != null) {
close();
pool = null;