This is an automated email from the ASF dual-hosted git repository.
jt2594838 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 60833764d31 Fix pipe historical extraction when device metadata is
unavailable (#18541)
60833764d31 is described below
commit 60833764d317cad23f9069e1cafdd52cdf5673e4
Author: Caideyipi <[email protected]>
AuthorDate: Mon Aug 31 14:58:48 2026 +0800
Fix pipe historical extraction when device metadata is unavailable (#18541)
---
...istoricalDataRegionTsFileAndDeletionSource.java | 2 +-
...ricalDataRegionTsFileAndDeletionSourceTest.java | 26 ++++++++++++++++++++++
2 files changed, 27 insertions(+), 1 deletion(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileAndDeletionSource.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileAndDeletionSource.java
index 1ee89144d88..e542fc51719 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileAndDeletionSource.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileAndDeletionSource.java
@@ -1012,7 +1012,7 @@ public class
PipeHistoricalDataRegionTsFileAndDeletionSource
.getDeviceIsAlignedMapFromCache(resource.getTsFile(), false);
deviceSet =
Objects.nonNull(deviceIsAlignedMap) ? deviceIsAlignedMap.keySet() :
resource.getDevices();
- } catch (final IOException e) {
+ } catch (final IOException | RuntimeException e) {
LOGGER.warn(
DataNodePipeMessages.PIPE_FAILED_TO_GET_DEVICES_FROM_TSFILE_1,
pipeName,
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileAndDeletionSourceTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileAndDeletionSourceTest.java
index 0dd6bae7ddf..991dab795f9 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileAndDeletionSourceTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileAndDeletionSourceTest.java
@@ -31,12 +31,14 @@ import
org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant;
import org.apache.iotdb.commons.pipe.config.constant.SystemConstant;
import
org.apache.iotdb.commons.pipe.config.plugin.configuraion.PipeTaskRuntimeConfiguration;
import
org.apache.iotdb.commons.pipe.config.plugin.env.PipeTaskSourceRuntimeEnvironment;
+import org.apache.iotdb.commons.pipe.datastructure.pattern.PrefixTreePattern;
import org.apache.iotdb.commons.pipe.datastructure.resource.PersistentResource;
import org.apache.iotdb.commons.pipe.event.ProgressReportEvent;
import org.apache.iotdb.commons.utils.FileUtils;
import org.apache.iotdb.db.pipe.consensus.ReplicateProgressDataNodeManager;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
import
org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResourceStatus;
+import
org.apache.iotdb.db.storageengine.dataregion.tsfile.timeindex.FileTimeIndex;
import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameterValidator;
import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameters;
import org.apache.iotdb.pipe.api.event.Event;
@@ -146,6 +148,30 @@ public class
PipeHistoricalDataRegionTsFileAndDeletionSourceTest {
}
}
+ @Test
+ public void testMissingTsFileResourceDoesNotBlockHistoricalExtraction()
throws Exception {
+ final PipeHistoricalDataRegionTsFileAndDeletionSource source =
+ new PipeHistoricalDataRegionTsFileAndDeletionSource();
+ final File tempDir =
Files.createTempDirectory("pipeHistoricalMissingResource").toFile();
+
+ try {
+ final TsFileResource resource = createTsFileResource(tempDir,
"missing-resource.tsfile");
+ resource.setTimeIndex(new FileTimeIndex());
+ setPrivateField(source, "pipeName", "pipe");
+ setPrivateField(source, "dataRegionId", 1);
+ setPrivateField(source, "treePattern", new PrefixTreePattern("root.**"));
+
+ final Method method =
+
PipeHistoricalDataRegionTsFileAndDeletionSource.class.getDeclaredMethod(
+ "mayTsFileResourceOverlappedWithPattern", TsFileResource.class);
+ method.setAccessible(true);
+
+ Assert.assertTrue((Boolean) method.invoke(source, resource));
+ } finally {
+ FileUtils.deleteFileOrDirectory(tempDir);
+ }
+ }
+
@Test
public void testSupplyRetriesSameTsFileAfterEventCreationFailure() throws
Exception {
final TestablePipeHistoricalDataRegionTsFileAndDeletionSource source =