From f1b6dc60bf8192696e3d4858549f958899da444e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=99=88=E5=93=B2=E6=B6=B5?= Date: Wed, 26 Aug 2026 03:19:24 +0000 Subject: [PATCH 1/2] [To dev/1.3] Load: preserve pending tablets when parser fails (#18487) (#18509) --- ...tementDataTypeConvertExecutionVisitor.java | 50 +++++++++++++++- ...ntDataTypeConvertExecutionVisitorTest.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 a1da70952461..3054588553f8 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.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 tabletIteratorFactory; @FunctionalInterface public interface StatementExecutor { @@ -66,7 +68,14 @@ public interface StatementExecutor { public LoadTreeStatementDataTypeConvertExecutionVisitor( final StatementExecutor statementExecutor) { + this(statementExecutor, file -> new LoadTreeTsFileTabletIterator(file, true)); + } + + LoadTreeStatementDataTypeConvertExecutionVisitor( + final StatementExecutor statementExecutor, + final Function tabletIteratorFactory) { this.statementExecutor = statementExecutor; + this.tabletIteratorFactory = tabletIteratorFactory; } @Override @@ -89,7 +98,7 @@ public Optional visitLoadFile( try { for (final File file : loadTsFileStatement.getTsFiles()) { try (final LoadTreeTsFileTabletIterator tabletIterator = - new LoadTreeTsFileTabletIterator(file, true)) { + tabletIteratorFactory.apply(file)) { for (final Pair tabletWithIsAligned : tabletIterator) { final PipeTransferTabletRawReq tabletRawReq = PipeTransferTabletRawReq.toTPipeTransferRawReq( @@ -123,9 +132,29 @@ public Optional visitLoadFile( } 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 Optional visitLoadFile( 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 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 150c950fce3e..0cfbbc364f4c 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.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 void testFallbackToQueryForRemainingDevicesWhenScanParserHitsCorruption() Assert.assertEquals(loadedPointCountBeforeCorruption, loadedPointCountAfterFallback); } + @Test + public void testFlushesPendingTabletsWhenIteratorFails() throws Exception { + tsFile = File.createTempFile("load-tree-pending-tablet", ".tsfile"); + final List 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 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 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 @@ private static void setPipeMemoryManagementEnabled(final boolean enabled) { CommonDescriptor.getInstance().getConfig().setPipeMemoryManagementEnabled(enabled); } } + + private static class ThrowingTabletIterator extends LoadTreeTsFileTabletIterator { + private final Pair 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 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. + } + } } From 07f173b43f54c11ea9dcb62bc6a3394b1215e34c Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Fri, 11 Sep 2026 12:49:19 +0800 Subject: [PATCH 2/2] Fix missing historical source test helper --- .../PipeHistoricalDataRegionTsFileSourceTest.java | 8 ++++++++ 1 file changed, 8 insertions(+) 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 5e68a17ba8e0..7ff080876924 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.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 @@ private static void assertMayTsFileContainUnprocessedData( 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 {