This is an automated email from the ASF dual-hosted git repository.
jt2594838 pushed a commit to branch rc/1.3.8
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/rc/1.3.8 by this push:
new f4f59e57718 [To dev/1.3] Load: preserve pending tablets when parser
fails (#18487) (#18509) (#18604)
f4f59e57718 is described below
commit f4f59e57718feda3b0fab91ce3222460c5798648
Author: Caideyipi <[email protected]>
AuthorDate: Fri Sep 18 12:37:18 2026 +0800
[To dev/1.3] Load: preserve pending tablets when parser fails (#18487)
(#18509) (#18604)
* [To dev/1.3] Load: preserve pending tablets when parser fails (#18487)
(#18509)
* Fix missing historical source test helper
---------
Co-authored-by: 陈哲涵 <[email protected]>
---
...eeStatementDataTypeConvertExecutionVisitor.java | 50 ++++++++++++++++--
.../PipeHistoricalDataRegionTsFileSourceTest.java | 8 +++
...atementDataTypeConvertExecutionVisitorTest.java | 60 ++++++++++++++++++++++
3 files changed, 115 insertions(+), 3 deletions(-)
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeStatementDataTypeConvertExecutionVisitor.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeStatementDataTypeConvertExecutionVisitor.java
index a1da7095246..3054588553f 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeStatementDataTypeConvertExecutionVisitor.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeStatementDataTypeConvertExecutionVisitor.java
@@ -43,6 +43,7 @@ import java.io.File;
import java.util.ArrayList;
import java.util.List;
import java.util.Optional;
+import java.util.function.Function;
import java.util.stream.Collectors;
import static
org.apache.iotdb.db.pipe.resource.memory.PipeMemoryWeightUtil.calculateTabletSizeInBytes;
@@ -58,6 +59,7 @@ public class LoadTreeStatementDataTypeConvertExecutionVisitor
.getLoadTsFileTabletConversionBatchMemorySizeInBytes();
private final StatementExecutor statementExecutor;
+ private final Function<File, LoadTreeTsFileTabletIterator>
tabletIteratorFactory;
@FunctionalInterface
public interface StatementExecutor {
@@ -66,7 +68,14 @@ public class LoadTreeStatementDataTypeConvertExecutionVisitor
public LoadTreeStatementDataTypeConvertExecutionVisitor(
final StatementExecutor statementExecutor) {
+ this(statementExecutor, file -> new LoadTreeTsFileTabletIterator(file,
true));
+ }
+
+ LoadTreeStatementDataTypeConvertExecutionVisitor(
+ final StatementExecutor statementExecutor,
+ final Function<File, LoadTreeTsFileTabletIterator>
tabletIteratorFactory) {
this.statementExecutor = statementExecutor;
+ this.tabletIteratorFactory = tabletIteratorFactory;
}
@Override
@@ -89,7 +98,7 @@ public class LoadTreeStatementDataTypeConvertExecutionVisitor
try {
for (final File file : loadTsFileStatement.getTsFiles()) {
try (final LoadTreeTsFileTabletIterator tabletIterator =
- new LoadTreeTsFileTabletIterator(file, true)) {
+ tabletIteratorFactory.apply(file)) {
for (final Pair<Tablet, Boolean> tabletWithIsAligned :
tabletIterator) {
final PipeTransferTabletRawReq tabletRawReq =
PipeTransferTabletRawReq.toTPipeTransferRawReq(
@@ -123,9 +132,29 @@ public class
LoadTreeStatementDataTypeConvertExecutionVisitor
} catch (final Exception e) {
LOGGER.warn(
"Failed to convert data type for LoadTsFileStatement: {}.",
loadTsFileStatement, e);
- return Optional.of(
+ final TSStatus status =
loadTsFileStatement.accept(
- LoadTsFileDataTypeConverter.STATEMENT_EXCEPTION_VISITOR, e));
+ LoadTsFileDataTypeConverter.STATEMENT_EXCEPTION_VISITOR, e);
+
+ // A parser can fail after producing tablets that are still waiting
for the next batch
+ // boundary. Submit those tablets before reporting the parser error
so the error does not
+ // discard successfully converted data.
+ if (!isRetryableConversionException(e) && !tabletRawReqs.isEmpty()) {
+ final TSStatus flushStatus =
+ executeInsertMultiTabletsWithRetry(
+ tabletRawReqs,
loadTsFileStatement.isConvertOnTypeMismatch());
+
+ for (final long memoryCost : tabletRawReqSizes) {
+ block.reduceMemoryUsage(memoryCost);
+ }
+ tabletRawReqs.clear();
+ tabletRawReqSizes.clear();
+
+ if (!handleTSStatus(flushStatus, loadTsFileStatement)) {
+ return Optional.of(flushStatus);
+ }
+ }
+ return Optional.of(status);
}
}
@@ -179,6 +208,21 @@ public class
LoadTreeStatementDataTypeConvertExecutionVisitor
return Optional.of(new
TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode()));
}
+ private static boolean isRetryableConversionException(final Throwable
throwable) {
+ if (LoadTsFileDataTypeConverter.isMemoryPressureException(throwable)) {
+ return true;
+ }
+
+ Throwable current = throwable;
+ while (current != null) {
+ if (current instanceof InterruptedException) {
+ return true;
+ }
+ current = current.getCause();
+ }
+ return false;
+ }
+
private TSStatus executeInsertMultiTabletsWithRetry(
final List<PipeTransferTabletRawReq> tabletRawReqs, boolean
isConvertOnTypeMismatch) {
final InsertMultiTabletsStatement batchStatement = new
InsertMultiTabletsStatement();
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileSourceTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileSourceTest.java
index 5e68a17ba8e..7ff08087692 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileSourceTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileSourceTest.java
@@ -40,6 +40,7 @@ import org.junit.Assert;
import org.junit.Test;
import java.io.File;
+import java.io.IOException;
import java.lang.reflect.Field;
import java.lang.reflect.Method;
import java.nio.file.Files;
@@ -198,6 +199,13 @@ public class PipeHistoricalDataRegionTsFileSourceTest {
source, createClosedTsFileResource(tempDir, fileName,
resourceProgressIndex)));
}
+ private static TsFileResource createTsFileResource(final File tempDir, final
String fileName)
+ throws IOException {
+ final File file = new File(tempDir, fileName);
+ Assert.assertTrue(file.createNewFile());
+ return new TsFileResource(file);
+ }
+
private static TsFileResource createClosedTsFileResource(
final File tempDir, final String fileName, final ProgressIndex
progressIndex)
throws Exception {
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeStatementDataTypeConvertExecutionVisitorTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeStatementDataTypeConvertExecutionVisitorTest.java
index 150c950fce3..0cfbbc364f4 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeStatementDataTypeConvertExecutionVisitorTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/load/converter/LoadTreeStatementDataTypeConvertExecutionVisitorTest.java
@@ -38,6 +38,7 @@ import org.apache.tsfile.file.metadata.enums.TSEncoding;
import org.apache.tsfile.read.TsFileSequenceReader;
import org.apache.tsfile.read.common.Path;
import org.apache.tsfile.utils.BitMap;
+import org.apache.tsfile.utils.Pair;
import org.apache.tsfile.write.TsFileWriter;
import org.apache.tsfile.write.record.Tablet;
import org.apache.tsfile.write.schema.MeasurementSchema;
@@ -110,6 +111,33 @@ public class
LoadTreeStatementDataTypeConvertExecutionVisitorTest {
Assert.assertEquals(loadedPointCountBeforeCorruption,
loadedPointCountAfterFallback);
}
+ @Test
+ public void testFlushesPendingTabletsWhenIteratorFails() throws Exception {
+ tsFile = File.createTempFile("load-tree-pending-tablet", ".tsfile");
+ final List<MeasurementSchema> schemaList =
+ Arrays.asList(new MeasurementSchema("s0", TSDataType.INT64,
TSEncoding.PLAIN));
+ final Tablet tablet = new Tablet(DEVICE_0, schemaList, 1);
+ tablet.addTimestamp(0, 1);
+ tablet.addValue("s0", 0, 1L);
+ tablet.rowSize = 1;
+
+ final Map<String, Integer> pointCountByDevice = new HashMap<>();
+ final LoadTreeStatementDataTypeConvertExecutionVisitor visitor =
+ new LoadTreeStatementDataTypeConvertExecutionVisitor(
+ statement -> {
+ collectLoadedPoints(statement, pointCountByDevice);
+ return new TSStatus(TSStatusCode.SUCCESS_STATUS.getStatusCode());
+ },
+ file -> new ThrowingTabletIterator(file, tablet));
+
+ final Optional<TSStatus> status =
+
visitor.visitLoadFile(LoadTsFileStatement.createUnchecked(tsFile.getAbsolutePath()),
null);
+
+ Assert.assertTrue(status.isPresent());
+ Assert.assertEquals(TSStatusCode.LOAD_FILE_ERROR.getStatusCode(),
status.get().getCode());
+ Assert.assertEquals(1, pointCountByDevice.getOrDefault(DEVICE_0,
0).intValue());
+ }
+
@Test
public void testFallbackToQueryWhenFirstNonAlignedDeviceIsCorrupted() throws
Exception {
tsFile = new
File("load-tree-query-fallback-corrupted-first-non-aligned-device.tsfile");
@@ -393,4 +421,36 @@ public class
LoadTreeStatementDataTypeConvertExecutionVisitorTest {
CommonDescriptor.getInstance().getConfig().setPipeMemoryManagementEnabled(enabled);
}
}
+
+ private static class ThrowingTabletIterator extends
LoadTreeTsFileTabletIterator {
+ private final Pair<Tablet, Boolean> tabletWithIsAligned;
+ private boolean tabletAvailable = true;
+
+ private ThrowingTabletIterator(final File file, final Tablet tablet) {
+ super(file, true);
+ tabletWithIsAligned = new Pair<>(tablet, false);
+ }
+
+ @Override
+ public boolean hasNext() {
+ if (tabletAvailable) {
+ return true;
+ }
+ throw new IllegalStateException("synthetic parser failure");
+ }
+
+ @Override
+ public Pair<Tablet, Boolean> next() {
+ if (!tabletAvailable) {
+ throw new IllegalStateException("synthetic parser failure");
+ }
+ tabletAvailable = false;
+ return tabletWithIsAligned;
+ }
+
+ @Override
+ public void close() {
+ // No parser resources are allocated by this test iterator.
+ }
+ }
}