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

rong 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 617bfcea3d3 Pipe: assign progress index for sealed tsfile recovery 
(#11384)
617bfcea3d3 is described below

commit 617bfcea3d3281fe84588717f57876b622515228
Author: Steve Yurong Su <[email protected]>
AuthorDate: Wed Oct 25 15:59:55 2023 +0800

    Pipe: assign progress index for sealed tsfile recovery (#11384)
---
 .../org/apache/iotdb/db/pipe/agent/runtime/PipeRuntimeAgent.java    | 6 +++---
 .../plan/planner/plan/node/load/LoadSingleTsFileNode.java           | 2 +-
 .../iotdb/db/storageengine/dataregion/tsfile/TsFileResource.java    | 2 +-
 .../dataregion/wal/recover/file/AbstractTsFileRecoverPerformer.java | 5 +++++
 .../dataregion/wal/recover/file/UnsealedTsFileRecoverPerformer.java | 2 +-
 .../src/main/java/org/apache/iotdb/db/utils/FileLoaderUtils.java    | 2 +-
 6 files changed, 12 insertions(+), 7 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/runtime/PipeRuntimeAgent.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/runtime/PipeRuntimeAgent.java
index dc50ba94bee..f7a2a5595f7 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/runtime/PipeRuntimeAgent.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/runtime/PipeRuntimeAgent.java
@@ -102,14 +102,14 @@ public class PipeRuntimeAgent implements IService {
 
   ////////////////////// Recover ProgressIndex Assigner //////////////////////
 
-  public void assignRecoverProgressIndexForTsFileRecovery(TsFileResource 
tsFileResource) {
-    tsFileResource.recoverProgressIndex(
+  public void assignProgressIndexForTsFileLoad(TsFileResource tsFileResource) {
+    tsFileResource.setProgressIndex(
         new RecoverProgressIndex(
             DATA_NODE_ID,
             
simpleConsensusProgressIndexAssigner.getSimpleProgressIndexForTsFileRecovery()));
   }
 
-  public void assignUpdateProgressIndexForTsFileRecovery(TsFileResource 
tsFileResource) {
+  public void assignProgressIndexForTsFileRecovery(TsFileResource 
tsFileResource) {
     tsFileResource.updateProgressIndex(
         new RecoverProgressIndex(
             DATA_NODE_ID,
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/load/LoadSingleTsFileNode.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/load/LoadSingleTsFileNode.java
index 687a652c363..897522afa04 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/load/LoadSingleTsFileNode.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/queryengine/plan/planner/plan/node/load/LoadSingleTsFileNode.java
@@ -100,7 +100,7 @@ public class LoadSingleTsFileNode extends WritePlanNode {
       needDecodeTsFile = !isDispatchedToLocal(new 
HashSet<>(partitionFetcher.apply(slotList)));
     }
 
-    PipeAgent.runtime().assignRecoverProgressIndexForTsFileRecovery(resource);
+    PipeAgent.runtime().assignProgressIndexForTsFileLoad(resource);
 
     // we serialize the resource file even if the tsfile does not need to be 
decoded
     // or the resource file is already existed because we need to serialize the
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/tsfile/TsFileResource.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/tsfile/TsFileResource.java
index ce49b700768..4aa7f546df3 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/tsfile/TsFileResource.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/tsfile/TsFileResource.java
@@ -1127,7 +1127,7 @@ public class TsFileResource {
             : 
maxProgressIndex.updateToMinimumIsAfterProgressIndex(progressIndex));
   }
 
-  public void recoverProgressIndex(ProgressIndex progressIndex) {
+  public void setProgressIndex(ProgressIndex progressIndex) {
     if (progressIndex == null) {
       return;
     }
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/recover/file/AbstractTsFileRecoverPerformer.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/recover/file/AbstractTsFileRecoverPerformer.java
index fbaba378da1..784a8c18045 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/recover/file/AbstractTsFileRecoverPerformer.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/recover/file/AbstractTsFileRecoverPerformer.java
@@ -20,6 +20,7 @@
 package org.apache.iotdb.db.storageengine.dataregion.wal.recover.file;
 
 import org.apache.iotdb.db.exception.DataRegionException;
+import org.apache.iotdb.db.pipe.agent.PipeAgent;
 import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
 import org.apache.iotdb.db.utils.FileLoaderUtils;
 import org.apache.iotdb.tsfile.exception.NotCompatibleTsFileException;
@@ -117,6 +118,10 @@ public abstract class AbstractTsFileRecoverPerformer 
implements Closeable {
         new 
TsFileSequenceReader(tsFileResource.getTsFile().getAbsolutePath())) {
       FileLoaderUtils.updateTsFileResource(reader, tsFileResource);
     }
+
+    // set progress index for pipe to avoid data loss
+    PipeAgent.runtime().assignProgressIndexForTsFileRecovery(tsFileResource);
+
     tsFileResource.serialize();
   }
 
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/recover/file/UnsealedTsFileRecoverPerformer.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/recover/file/UnsealedTsFileRecoverPerformer.java
index 9c883889143..b0aaa8b1378 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/recover/file/UnsealedTsFileRecoverPerformer.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/wal/recover/file/UnsealedTsFileRecoverPerformer.java
@@ -243,7 +243,7 @@ public class UnsealedTsFileRecoverPerformer extends 
AbstractTsFileRecoverPerform
         }
 
         // set recover progress index for pipe
-        
PipeAgent.runtime().assignUpdateProgressIndexForTsFileRecovery(tsFileResource);
+        
PipeAgent.runtime().assignProgressIndexForTsFileRecovery(tsFileResource);
 
         // if we put following codes in the 'if' clause above, this file can 
be continued writing
         // into it
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/FileLoaderUtils.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/FileLoaderUtils.java
index 2d6d53809da..eccb532e8df 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/FileLoaderUtils.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/utils/FileLoaderUtils.java
@@ -108,7 +108,7 @@ public class FileLoaderUtils {
       }
     }
     resource.setStatus(TsFileResourceStatus.NORMAL);
-    PipeAgent.runtime().assignRecoverProgressIndexForTsFileRecovery(resource);
+    PipeAgent.runtime().assignProgressIndexForTsFileLoad(resource);
     return resource;
   }
 

Reply via email to