This is an automated email from the ASF dual-hosted git repository. Caideyipi pushed a commit to branch to-138-gr in repository https://gitbox.apache.org/repos/asf/iotdb.git
commit f1b6dc60bf8192696e3d4858549f958899da444e Author: 陈哲涵 <[email protected]> AuthorDate: Wed Aug 26 03:19:24 2026 +0000 [To dev/1.3] Load: preserve pending tablets when parser fails (#18487) (#18509) --- ...eeStatementDataTypeConvertExecutionVisitor.java | 50 ++++++++++++++++-- ...atementDataTypeConvertExecutionVisitorTest.java | 60 ++++++++++++++++++++++ 2 files changed, 107 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/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. + } + } }
