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;
}