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