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

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


The following commit(s) were added to refs/heads/master by this push:
     new e6713d4f029 IoTV2: Clean receiver files when dropping consensus pipe & 
Improve robustness when cleaning some dirs. (#15252)
e6713d4f029 is described below

commit e6713d4f029f0ab416cade9ca735b06cc4ea2e72
Author: Peng Junzhi <[email protected]>
AuthorDate: Wed Apr 2 16:53:56 2025 +0800

    IoTV2: Clean receiver files when dropping consensus pipe & Improve 
robustness when cleaning some dirs. (#15252)
    
    * test conf
    
    * clean receiver files after drop pipes
    
    * Revert "test conf"
    
    This reverts commit a839dfe3335c49645d0ad64f0b3a3d684783ca60.
---
 .../iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java   | 14 ++++++++++++++
 .../pipe/consensus/ConsensusPipeDataNodeDispatcher.java   |  3 ---
 .../pipe/consensus/deletion/DeletionResourceManager.java  | 15 ++++++++++++---
 .../pipeconsensus/PipeConsensusReceiverAgent.java         |  2 +-
 .../thrift/impl/DataNodeInternalRPCServiceImpl.java       |  4 ----
 5 files changed, 27 insertions(+), 11 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java
index 73606a04252..b6f3cffce33 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java
@@ -331,6 +331,13 @@ public class PipeDataNodeTaskAgent extends PipeTaskAgent {
     PipeTsFileToTabletsMetrics.getInstance().deregister(taskId);
     PipeDataNodeRemainingEventAndTimeMetrics.getInstance().deregister(taskId);
 
+    if (pipeName.startsWith(PipeStaticMeta.CONSENSUS_PIPE_PREFIX)) {
+      // Release corresponding receiver's resource
+      PipeDataNodeAgent.receiver()
+          .pipeConsensus()
+          .handleDropPipeConsensusTask(new ConsensusPipeName(pipeName));
+    }
+
     return true;
   }
 
@@ -350,6 +357,13 @@ public class PipeDataNodeTaskAgent extends PipeTaskAgent {
       
PipeDataNodeRemainingEventAndTimeMetrics.getInstance().deregister(taskId);
     }
 
+    if (pipeName.startsWith(PipeStaticMeta.CONSENSUS_PIPE_PREFIX)) {
+      // Release corresponding receiver's resource
+      PipeDataNodeAgent.receiver()
+          .pipeConsensus()
+          .handleDropPipeConsensusTask(new ConsensusPipeName(pipeName));
+    }
+
     return true;
   }
 
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/consensus/ConsensusPipeDataNodeDispatcher.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/consensus/ConsensusPipeDataNodeDispatcher.java
index 3e56a57e3ca..02179a29f56 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/consensus/ConsensusPipeDataNodeDispatcher.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/consensus/ConsensusPipeDataNodeDispatcher.java
@@ -25,7 +25,6 @@ import org.apache.iotdb.commons.consensus.ConfigRegionId;
 import org.apache.iotdb.confignode.rpc.thrift.TCreatePipeReq;
 import org.apache.iotdb.consensus.pipe.consensuspipe.ConsensusPipeDispatcher;
 import org.apache.iotdb.consensus.pipe.consensuspipe.ConsensusPipeName;
-import org.apache.iotdb.db.pipe.agent.PipeDataNodeAgent;
 import org.apache.iotdb.db.protocol.client.ConfigNodeClient;
 import org.apache.iotdb.db.protocol.client.ConfigNodeClientManager;
 import org.apache.iotdb.db.protocol.client.ConfigNodeInfo;
@@ -128,7 +127,5 @@ public class ConsensusPipeDataNodeDispatcher implements 
ConsensusPipeDispatcher
       LOGGER.warn("Failed to drop consensus pipe-{}", pipeName, e);
       throw new PipeException("Failed to drop consensus pipe", e);
     }
-    // Release corresponding receiver's resource
-    
PipeDataNodeAgent.receiver().pipeConsensus().handleDropPipeConsensusTask(pipeName);
   }
 }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/consensus/deletion/DeletionResourceManager.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/consensus/deletion/DeletionResourceManager.java
index dcde242460c..ff868ffd445 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/consensus/deletion/DeletionResourceManager.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/consensus/deletion/DeletionResourceManager.java
@@ -168,9 +168,18 @@ public class DeletionResourceManager implements 
AutoCloseable {
   }
 
   public void removeDAL() {
-    FileUtils.deleteFileOrDirectory(storageDir);
-    LOGGER.info(
-        "DeletionManager-{}: current DAL dir {} is deleted successfully", 
dataRegionId, storageDir);
+    if (storageDir.exists()) {
+      FileUtils.deleteFileOrDirectory(storageDir);
+      LOGGER.info(
+          "DeletionManager-{}: current DAL dir {} is deleted successfully",
+          dataRegionId,
+          storageDir);
+    } else {
+      LOGGER.info(
+          "DeletionManager-{}: current DAL dir {} is not initialized, no need 
to delete.",
+          dataRegionId,
+          storageDir);
+    }
   }
 
   /**
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/pipeconsensus/PipeConsensusReceiverAgent.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/pipeconsensus/PipeConsensusReceiverAgent.java
index ff3307c5f56..524b93dda4c 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/pipeconsensus/PipeConsensusReceiverAgent.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/receiver/protocol/pipeconsensus/PipeConsensusReceiverAgent.java
@@ -194,7 +194,7 @@ public class PipeConsensusReceiverAgent implements 
ConsensusPipeReceiver {
     if (receiverReference != null) {
       receiverReference.get().handleExit();
       receiverReference.set(null);
+      consensusPipe2ReciverMap.remove(pipeName);
     }
-    consensusPipe2ReciverMap.remove(pipeName);
   }
 }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImpl.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImpl.java
index 47fa540f51b..82e85276e64 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImpl.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImpl.java
@@ -82,7 +82,6 @@ import org.apache.iotdb.db.consensus.DataRegionConsensusImpl;
 import org.apache.iotdb.db.consensus.SchemaRegionConsensusImpl;
 import org.apache.iotdb.db.exception.StorageEngineException;
 import org.apache.iotdb.db.pipe.agent.PipeDataNodeAgent;
-import org.apache.iotdb.db.pipe.consensus.deletion.DeletionResourceManager;
 import org.apache.iotdb.db.protocol.client.ConfigNodeInfo;
 import 
org.apache.iotdb.db.protocol.client.cn.DnToCnInternalServiceAsyncRequestManager;
 import org.apache.iotdb.db.protocol.client.cn.DnToCnRequestType;
@@ -2306,9 +2305,6 @@ public class DataNodeInternalRPCServiceImpl implements 
IDataNodeRPCService.Iface
     if (consensusGroupId instanceof DataRegionId) {
       try {
         
DataRegionConsensusImpl.getInstance().deleteLocalPeer(consensusGroupId);
-        Optional.ofNullable(
-                
DeletionResourceManager.getInstance(String.valueOf(tconsensusGroupId.getId())))
-            .ifPresent(DeletionResourceManager::close);
       } catch (ConsensusException e) {
         if (!(e instanceof ConsensusGroupNotExistException)) {
           return RpcUtils.getStatus(TSStatusCode.DELETE_REGION_ERROR, 
e.getMessage());

Reply via email to